返回博客

Flight SQL:高性能数据传输协议

Tags: #FlightSQL #ArrowFlight #HighPerformance #DataTransfer #JDBC #智策平台

Coomia发布于 2025年8月3日9 分钟阅读
分享本文Twitter / X

系列:S3 数据基座 · 第 24 篇 | 难度:高级 | 阅读时间:20 分钟

Flight SQL:高性能数据传输协议

Tags: #FlightSQL #ArrowFlight #HighPerformance #DataTransfer #JDBC #智策平台

#TL;DR

Arrow Flight SQL 是 coomia-dip 平台高性能数据传输的核心协议。相比传统 JDBC/ODBC 的行序列化传输,Flight SQL 基于 Apache Arrow 列式内存格式和 gRPC 传输,实现了零拷贝数据交换,吞吐量提升 10-20 倍。本文完整解析 Flight SQL 在 coomia-dip 中的应用,包括协议架构、Doris Flight SQL 集成、Python SDK 客户端实现、OQL 查询通过 Flight SQL 执行的全链路、流式结果集处理、连接池管理、安全认证(mTLS + Token)和性能调优。

#1. 为什么选择 Flight SQL

#1.1 传统协议的瓶颈

Code
JDBC/ODBC vs Arrow Flight SQL:

传统 JDBC/ODBC 数据传输:
  Server → 列式存储 → 行序列化 → 网络传输 → 行反序列化 → 客户端
  瓶颈:序列化/反序列化占 CPU 的 70-90%

Arrow Flight SQL 数据传输:
  Server → Arrow 列式格式 → 网络传输 → 客户端(零拷贝)
  优势:数据已是列式格式,无需序列化转换

性能对比(传输 1000 万行、20 列):

┌─────────────────┬──────────┬──────────┬──────────┐
│ 协议              │ 吞吐量    │ 延迟      │ CPU 占用  │
├─────────────────┼──────────┼──────────┼──────────┤
│ JDBC (MySQL)     │ 80 MB/s  │ 高        │ 高        │
│ JDBC (Doris)     │ 120 MB/s │ 中        │ 中        │
│ Arrow Flight     │ 2 GB/s   │ 低        │ 低        │
│ Arrow Flight SQL │ 1.5 GB/s │ 低        │ 低        │
└─────────────────┴──────────┴──────────┴──────────┘

#1.2 Flight SQL 协议栈

Code
Arrow Flight SQL 协议栈:

┌──────────────────────────────────────┐
│          Application Layer            │
│  (OQL Query / SQL Query)              │
├──────────────────────────────────────┤
│          Flight SQL Layer             │
│  (GetFlightInfo / DoGet / DoPut)      │
├──────────────────────────────────────┤
│          Arrow Flight Layer           │
│  (Arrow IPC Format / Streaming)       │
├──────────────────────────────────────┤
│          gRPC Transport               │
│  (HTTP/2, multiplexing, flow control) │
├──────────────────────────────────────┤
│          TLS / mTLS                   │
│  (加密传输、双向认证)                   │
└──────────────────────────────────────┘

#2. Doris Flight SQL 集成

#2.1 Doris 端配置

Code
Doris Flight SQL 服务配置:

fe.conf:
  arrow_flight_sql_port = 8040
  arrow_flight_token_alive_time = 600

be.conf:
  arrow_flight_sql_port = 8060
  arrow_flight_result_batch_size = 8192

架构:
  Client → FE (8040) → 获取查询计划和 Ticket
         → BE (8060) → 流式获取数据(直连 BE 节点)

优势:
  客户端直连数据所在的 BE 节点,避免 FE 数据中转
  多 BE 并行传输,水平扩展吞吐量

#2.2 查询执行流程

Code
Flight SQL 查询执行流程:

1. Client 连接 FE
   FlightClient.connect("grpc://doris-fe:8040")
   FlightClient.authenticate(user, password)

2. Client 发送查询
   FlightInfo info = client.getFlightInfo(
     FlightDescriptor.command("SELECT * FROM entity_common LIMIT 1000")
   )

3. FE 返回 FlightInfo
   FlightInfo {
     schema: Arrow Schema,
     endpoints: [
       {ticket: "query-123-be1", locations: ["grpc://be1:8060"]},
       {ticket: "query-123-be2", locations: ["grpc://be2:8060"]},
       {ticket: "query-123-be3", locations: ["grpc://be3:8060"]},
     ]
   }

4. Client 并行从各 BE 获取数据
   for endpoint in info.endpoints:
       be_client = FlightClient.connect(endpoint.locations[0])
       stream = be_client.doGet(endpoint.ticket)
       for batch in stream:
           process(batch)  # Arrow RecordBatch, 零拷贝

#3. Python SDK 客户端

#3.1 Flight SQL 客户端封装

Python
import pyarrow.flight as flight
import pyarrow as pa

class OntoFlightSQLClient:
    """coomia-dip Flight SQL 客户端"""

    def __init__(self, host: str, port: int = 8040):
        self._location = flight.Location.for_grpc_tcp(host, port)
        self._client = flight.FlightClient(self._location)
        self._token: bytes | None = None

    def authenticate(self, username: str, password: str):
        """认证"""
        auth_handler = flight.ClientAuthHandler()
        self._token = self._client.authenticate_basic_token(
            username, password
        )

    def execute_query(self, sql: str) -> pa.Table:
        """执行查询并返回完整结果"""
        options = flight.FlightCallOptions(
            headers=[(b"authorization", self._token)]
        )

        # 获取 FlightInfo
        info = self._client.get_flight_info(
            flight.FlightDescriptor.for_command(sql.encode()),
            options=options
        )

        # 从所有 endpoint 并行获取数据
        batches = []
        for endpoint in info.endpoints:
            reader = self._client.do_get(
                endpoint.ticket, options=options
            )
            for batch in reader:
                batches.append(batch.data)

        if not batches:
            return pa.table({})

        return pa.Table.from_batches(batches, schema=info.schema)

    async def execute_streaming(
        self, sql: str
    ) -> AsyncIterator[pa.RecordBatch]:
        """流式执行查询"""
        options = flight.FlightCallOptions(
            headers=[(b"authorization", self._token)]
        )

        info = self._client.get_flight_info(
            flight.FlightDescriptor.for_command(sql.encode()),
            options=options
        )

        for endpoint in info.endpoints:
            reader = self._client.do_get(
                endpoint.ticket, options=options
            )
            for batch in reader:
                yield batch.data

#3.2 OQL 通过 Flight SQL 执行

Python
class OQLFlightExecutor:
    """OQL 查询通过 Flight SQL 执行"""

    def __init__(self, flight_client: OntoFlightSQLClient):
        self._client = flight_client
        self._compiler = OQLToSQLCompiler()

    def execute_oql(self, oql: str) -> pa.Table:
        # 步骤 1:OQL → SQL
        sql = self._compiler.compile(oql)

        # 步骤 2:通过 Flight SQL 执行
        result = self._client.execute_query(sql)

        # 步骤 3:后处理(指标展开等)
        result = self._post_process(result, oql)

        return result

    def execute_oql_to_pandas(self, oql: str) -> 'pd.DataFrame':
        """执行 OQL 并返回 Pandas DataFrame"""
        table = self.execute_oql(oql)
        return table.to_pandas()

    def execute_oql_to_polars(self, oql: str) -> 'pl.DataFrame':
        """执行 OQL 并返回 Polars DataFrame"""
        table = self.execute_oql(oql)
        return pl.from_arrow(table)

#4. 连接池管理

Python
class FlightConnectionPool:
    """Flight SQL 连接池"""

    def __init__(
        self, host: str, port: int,
        min_connections: int = 5,
        max_connections: int = 20,
        idle_timeout: int = 300,
    ):
        self._host = host
        self._port = port
        self._min = min_connections
        self._max = max_connections
        self._idle_timeout = idle_timeout
        self._pool: asyncio.Queue[OntoFlightSQLClient] = asyncio.Queue()
        self._active_count = 0

    async def acquire(self) -> OntoFlightSQLClient:
        try:
            client = self._pool.get_nowait()
            if client.is_healthy():
                return client
            self._active_count -= 1
        except asyncio.QueueEmpty:
            pass

        if self._active_count < self._max:
            client = OntoFlightSQLClient(self._host, self._port)
            client.authenticate(self._username, self._password)
            self._active_count += 1
            return client

        return await self._pool.get()

    async def release(self, client: OntoFlightSQLClient):
        if self._pool.qsize() < self._min:
            await self._pool.put(client)
        else:
            client.close()
            self._active_count -= 1

#5. 安全认证

Code
Flight SQL 安全认证方案:

方案 1:Basic Auth + Token
  1. 客户端发送 username/password
  2. 服务端返回 Bearer Token
  3. 后续请求携带 Token
  适用:内部服务间通信

方案 2:mTLS(双向 TLS)
  1. 客户端和服务端都有证书
  2. 建立连接时双向验证
  3. 传输层加密
  适用:跨网络、高安全要求场景

方案 3:JWT Token
  1. 通过 OAuth2/OIDC 获取 JWT
  2. Flight SQL 请求携带 JWT
  3. 服务端验证 JWT 签名和权限
  适用:与企业 SSO 集成

#6. 性能调优

Code
Flight SQL 性能调优参数:

┌──────────────────────┬──────────┬──────────────────┐
│ 参数                    │ 默认值    │ 推荐值             │
├──────────────────────┼──────────┼──────────────────┤
│ batch_size             │ 4096     │ 8192-16384       │
│ max_message_size       │ 4 MB     │ 16 MB            │
│ grpc_keepalive_time    │ 120s     │ 30s              │
│ flight_parallelism     │ 1        │ BE 节点数         │
│ compression            │ none     │ lz4 / zstd       │
│ connection_pool_size   │ 5        │ 10-20            │
│ timeout_seconds        │ 30       │ 60-300           │
└──────────────────────┴──────────┴──────────────────┘

批量大小调优:
  batch_size 太小 → gRPC 开销大
  batch_size 太大 → 单批延迟高
  最优:8192-16384(根据列数和列宽调整)

压缩调优:
  无压缩:最低 CPU,最高带宽
  LZ4:低 CPU,中等压缩(推荐)
  ZSTD:中 CPU,高压缩(带宽受限时)

#7. 测试策略

Python
class TestFlightSQL:

    def test_basic_query(self):
        client = OntoFlightSQLClient("localhost", 8040)
        client.authenticate("admin", "password")
        result = client.execute_query("SELECT 1 AS n")
        assert result.num_rows == 1
        assert result.column('n')[0].as_py() == 1

    def test_large_result_set(self):
        client = OntoFlightSQLClient("localhost", 8040)
        client.authenticate("admin", "password")
        result = client.execute_query(
            "SELECT * FROM entity_common LIMIT 1000000"
        )
        assert result.num_rows == 1000000

    def test_parallel_endpoint_fetch(self):
        client = OntoFlightSQLClient("localhost", 8040)
        client.authenticate("admin", "password")
        info = client.get_flight_info(
            "SELECT * FROM entity_common"
        )
        # 应有多个 endpoint(多 BE 节点)
        assert len(info.endpoints) >= 1

    def test_oql_via_flight_sql(self):
        executor = OQLFlightExecutor(client)
        result = executor.execute_oql(
            "FETCH Person WHERE age > 30 LIMIT 100"
        )
        assert result.num_rows <= 100

    async def test_connection_pool(self):
        pool = FlightConnectionPool("localhost", 8040)
        # 并发获取 10 个连接
        clients = await asyncio.gather(*[
            pool.acquire() for _ in range(10)
        ])
        assert len(clients) == 10
        for c in clients:
            await pool.release(c)

    def test_throughput_benchmark(self):
        start = time.time()
        result = client.execute_query(
            "SELECT * FROM entity_common LIMIT 10000000"
        )
        elapsed = time.time() - start
        throughput = result.nbytes / elapsed / 1024 / 1024
        assert throughput > 500  # > 500 MB/s

#Key Takeaways

  1. Arrow Flight SQL 将数据传输吞吐量提升 10-20 倍:基于 Arrow 列式格式的零拷贝传输,消除了传统 JDBC 的序列化/反序列化瓶颈。

  2. 直连 BE 节点实现并行数据获取:查询计划返回多个 endpoint,客户端直连数据所在的 BE 节点并行读取,吞吐量随 BE 节点数线性扩展。

  3. 与 Python 数据科学生态无缝集成:Arrow Table 可以零拷贝转换为 Pandas DataFrame 或 Polars DataFrame,适合数据分析和机器学习场景。

  4. 连接池和压缩是生产环境必须:连接池避免频繁建连开销,LZ4 压缩在几乎不增加 CPU 的情况下减少 50-70% 的网络传输。

  5. OQL 通过 Flight SQL 执行实现了端到端高性能:OQL → SQL → Flight SQL → Arrow Table 的全链路避免了任何中间格式转换。

#Next Article

下一篇 S3-25《时序数据:IoT 与监控场景的数据基座》 是 S3 系列的收官之作,将展示如何在 Ontology 三表模型上支持时序数据场景。

Tags: #FlightSQL #ArrowFlight #HighPerformance #DataTransfer #ZeroCopy #gRPC #ConnectionPool #智策平台 #coomia-dip #数据基座