Flight SQL:高性能数据传输协议
Tags: #FlightSQL #ArrowFlight #HighPerformance #DataTransfer #JDBC #智策平台
“系列: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 传统协议的瓶颈
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 协议栈
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 端配置
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 查询执行流程
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 客户端封装
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 执行
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. 连接池管理
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. 安全认证
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. 性能调优
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. 测试策略
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
-
Arrow Flight SQL 将数据传输吞吐量提升 10-20 倍:基于 Arrow 列式格式的零拷贝传输,消除了传统 JDBC 的序列化/反序列化瓶颈。
-
直连 BE 节点实现并行数据获取:查询计划返回多个 endpoint,客户端直连数据所在的 BE 节点并行读取,吞吐量随 BE 节点数线性扩展。
-
与 Python 数据科学生态无缝集成:Arrow Table 可以零拷贝转换为 Pandas DataFrame 或 Polars DataFrame,适合数据分析和机器学习场景。
-
连接池和压缩是生产环境必须:连接池避免频繁建连开销,LZ4 压缩在几乎不增加 CPU 的情况下减少 50-70% 的网络传输。
-
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 #数据基座