Arrow Flight SQL 深潜:高性能柱状数据传输
1. [为什么需要 Arrow Flight SQL](#1-为什么需要-arrow-flight-sql)
Coomia发布于 2025年11月25日10 分钟阅读
分享本文Twitter / X
“系列:S8 技术组件深潜 · 第 16 篇 | 难度:高级 | 阅读时间:20 分钟
Arrow Flight SQL 深潜:高性能柱状数据传输
#TL;DR
- Apache Arrow Flight SQL 是 coomia-dip 实现高性能数据传输的关键技术,将查询结果以零拷贝柱状格式在服务间传递,比 JSON 序列化快 10-100 倍
- 本文深入分析 Arrow 内存格式、Flight RPC 协议、Flight SQL 扩展、以及 coomia-dip 中 Data Layer 到 SDK 的高速数据通道
- 涵盖 Arrow IPC 格式、RecordBatch 流式传输、DoPut/DoGet 协议、与 Doris/Trino 的集成,以及内存管理与背压策略
#目录
- 为什么需要 Arrow Flight SQL
- Arrow 内存格式
- Flight RPC 协议
- Flight SQL 扩展
- coomia-dip 的数据传输架构
- Python 客户端实现
- Java 服务端实现
- 与 Doris 的集成
- 与 Trino 的集成
- 性能基准与优化
- Key Takeaways
#1. 为什么需要 Arrow Flight SQL
#1.1 传统数据传输的瓶颈
Code
传统 JSON/REST 数据传输:
Client ←─ JSON ─ Server
│
├── 序列化开销:Object → JSON String(CPU 密集)
├── 网络开销:JSON 文本体积是二进制的 3-5 倍
├── 反序列化开销:JSON String → Object(CPU 密集)
└── 内存分配:每行一个对象 → GC 压力大
Arrow Flight SQL 数据传输:
Client ←─ Arrow IPC ─ Server
│
├── 零序列化:内存格式即传输格式
├── 网络高效:二进制柱状格式,紧凑
├── 零拷贝:直接映射内存
└── SIMD 友好:柱状排列利于向量化计算
#1.2 性能对比
| 场景 | JSON/REST | Arrow Flight SQL | 加速比 |
|---|---|---|---|
| 1M 行 × 10 列(数值) | 12.5 秒 | 0.3 秒 | 42x |
| 1M 行 × 10 列(混合) | 18.2 秒 | 0.8 秒 | 23x |
| 100K 行 × 50 列 | 8.7 秒 | 0.2 秒 | 44x |
| 流式传输 10M 行 | OOM | 4.2 秒 | - |
#2. Arrow 内存格式
#2.1 柱状布局
Code
行式存储 (JSON/JDBC):
Row 0: {id: 1, name: "Alice", score: 95.5}
Row 1: {id: 2, name: "Bob", score: 87.3}
Row 2: {id: 3, name: "Carol", score: 92.1}
柱状存储 (Arrow):
Column "id": [1, 2, 3] → Int64 Buffer
Column "name": ["Alice", "Bob", "Carol"] → Utf8 Buffer (offset + data)
Column "score": [95.5, 87.3, 92.1] → Float64 Buffer
每列连续存储 → 缓存友好 → SIMD 向量化
#2.2 Buffer 结构
Code
Int64 Column "id" (3 values):
┌──────────────┬────────────────────────────┐
│ Validity Bitmap │ Data Buffer (24 bytes) │
│ [1, 1, 1] │ [01 00 00 00 00 00 00 00│
│ (3 bits → 1B) │ 02 00 00 00 00 00 00 00│
│ │ 03 00 00 00 00 00 00 00]│
└──────────────┴────────────────────────────┘
Utf8 Column "name" (3 values):
┌──────────────┬──────────────────┬──────────────────────┐
│ Validity │ Offset Buffer │ Data Buffer │
│ [1, 1, 1] │ [0, 5, 8, 13] │ "AliceBobCarol" │
└──────────────┴──────────────────┴──────────────────────┘
#2.3 RecordBatch 与 Schema
Python
import pyarrow as pa
# 定义 Schema
schema = pa.schema([
pa.field("id", pa.int64(), nullable=False),
pa.field("object_type", pa.utf8()),
pa.field("properties", pa.map_(pa.utf8(), pa.utf8())),
pa.field("created_at", pa.timestamp("us", tz="UTC")),
pa.field("version", pa.int32()),
])
# 创建 RecordBatch
batch = pa.record_batch(
[
pa.array([1, 2, 3], type=pa.int64()),
pa.array(["Employee", "Employee", "Department"]),
pa.array([{"name": "Alice"}, {"name": "Bob"}, {"name": "Engineering"}],
type=pa.map_(pa.utf8(), pa.utf8())),
pa.array([datetime.now()] * 3, type=pa.timestamp("us", tz="UTC")),
pa.array([1, 1, 1], type=pa.int32()),
],
schema=schema,
)
print(f"Batch size: {batch.nbytes} bytes, rows: {batch.num_rows}")
#3. Flight RPC 协议
#3.1 协议概览
Code
Flight 协议基于 gRPC,定义了 4 个核心 RPC:
1. GetFlightInfo(FlightDescriptor) → FlightInfo
"这个数据集有多大?在哪里?"
2. DoGet(Ticket) → Stream<FlightData>
"给我数据"(服务端流式传输)
3. DoPut(Stream<FlightData>) → PutResult
"接收数据"(客户端流式传输)
4. DoExchange(Stream<FlightData>) → Stream<FlightData>
"双向数据交换"(双向流)
#3.2 数据流
Code
Client Flight Server
│ │
│── GetFlightInfo(descriptor) ────→│
│←── FlightInfo(endpoints, size) ──│
│ │
│── DoGet(ticket) ────────────────→│
│←── Schema ───────────────────────│
│←── RecordBatch 1 ───────────────│
│←── RecordBatch 2 ───────────────│
│←── RecordBatch N ───────────────│
│←── (stream end) ────────────────│
#3.3 并行传输
Code
GetFlightInfo
│
FlightInfo {
endpoints: [
{ticket: T1, locations: [server-1]},
{ticket: T2, locations: [server-2]},
{ticket: T3, locations: [server-3]},
]
}
│
┌──────────┼──────────┐
│ │ │
DoGet(T1) DoGet(T2) DoGet(T3)
server-1 server-2 server-3
│ │ │
└──────────┼──────────┘
Client 并行接收
#4. Flight SQL 扩展
#4.1 Flight SQL 在 Flight 之上的扩展
Code
Flight SQL = Flight + SQL 语义
核心 RPC 扩展:
GetFlightInfo + CommandStatementQuery → 执行 SQL 查询
DoGet → 获取查询结果
DoPut + CommandStatementUpdate → 执行 DML
GetCatalogs / GetSchemas / GetTables → 元数据查询
#4.2 查询执行流程
Python
from adbc_driver_flightsql import dbapi as flight_sql
# 连接 Flight SQL 服务器
conn = flight_sql.connect(
"grpc://data-Layer:50051",
db_kwargs={
"username": "coomia-dip",
"password": "secret",
"adbc.flight.sql.client_option.tls_skip_verify": "true",
},
)
cursor = conn.cursor()
# 执行查询(返回 Arrow RecordBatch)
cursor.execute("""
SELECT id, object_type, properties, created_at
FROM ontology_instances
WHERE world_id = ? AND object_type = ?
""", ["world-001", "Employee"])
# 获取 Arrow Table(零拷贝)
table = cursor.fetch_arrow_table()
print(f"Rows: {table.num_rows}, Size: {table.nbytes / 1024:.1f} KB")
# 转换为 Pandas(零拷贝如果类型兼容)
df = table.to_pandas()
#5. coomia-dip 的数据传输架构
#5.1 架构图
Code
┌─────────────────────────────────────────────────┐
│ Python SDK │
│ ┌────────────────────────────────────────┐ │
│ │ ADBC Flight SQL Client │ │
│ │ → Arrow RecordBatch → Pandas / Polars │ │
│ └───────────────┬────────────────────────┘ │
└──────────────────┼──────────────────────────────┘
│ Flight SQL (gRPC)
┌──────────────────┼──────────────────────────────┐
│ Data Layer │ │
│ ┌───────────────┴────────────────────┐ │
│ │ Flight SQL Server (Quarkus) │ │
│ │ → SQL 解析 → 路由到存储引擎 │ │
│ └───────┬──────────────┬─────────────┘ │
│ │ │ │
│ ┌───────┴───┐ ┌──────┴──────┐ │
│ │ Doris │ │ Trino │ │
│ │ (OLAP) │ │ (Federation)│ │
│ └───────────┘ └─────────────┘ │
└─────────────────────────────────────────────────┘
#5.2 OQL 到 Flight SQL 的转换
Python
class OqlToFlightSqlTranslator:
"""将 OQL 查询转换为 Flight SQL 可执行的 SQL"""
def translate(self, oql: str, world_id: str) -> str:
parsed = self.parser.parse(oql)
if parsed.type == "FIND":
return self._translate_find(parsed, world_id)
elif parsed.type == "AGGREGATE":
return self._translate_aggregate(parsed, world_id)
def _translate_find(self, parsed, world_id: str) -> str:
object_type = parsed.object_type
table = f"ontology.{object_type.lower()}"
sql_parts = [f"SELECT * FROM {table}"]
sql_parts.append(f"WHERE world_id = '{world_id}'")
if parsed.filters:
for f in parsed.filters:
sql_parts.append(f"AND properties->>'{f.field}' {f.op} '{f.value}'")
if parsed.order_by:
sql_parts.append(f"ORDER BY {parsed.order_by}")
if parsed.limit:
sql_parts.append(f"LIMIT {parsed.limit}")
return " ".join(sql_parts)
#6. Python 客户端实现
#6.1 SDK 中的 Flight SQL 客户端
Python
import pyarrow.flight as flight
class CoomiaDipFlightClient:
"""coomia-dip Flight SQL 高性能数据客户端"""
def __init__(self, host: str, port: int = 50051):
self.client = flight.FlightClient(
f"grpc://{host}:{port}",
middleware=[AuthMiddleware()],
)
def query(self, sql: str, world_id: str) -> pa.Table:
"""执行查询并返回 Arrow Table"""
descriptor = flight.FlightDescriptor.for_command(
sql.encode("utf-8")
)
flight_info = self.client.get_flight_info(
descriptor,
options=flight.FlightCallOptions(
headers=[(b"x-world-id", world_id.encode())],
),
)
# 并行获取所有 Endpoint 的数据
batches = []
for endpoint in flight_info.endpoints:
reader = self.client.do_get(endpoint.ticket)
batches.append(reader.read_all())
return pa.concat_tables(batches)
def stream_query(self, sql: str, world_id: str):
"""流式查询,逐批返回 RecordBatch"""
descriptor = flight.FlightDescriptor.for_command(sql.encode())
flight_info = self.client.get_flight_info(descriptor)
for endpoint in flight_info.endpoints:
reader = self.client.do_get(endpoint.ticket)
for batch in reader:
yield batch.data # RecordBatch
#6.2 与 Pandas/Polars 集成
Python
# Pandas 集成
def query_to_pandas(client: CoomiaDipFlightClient, oql: str, world_id: str) -> pd.DataFrame:
table = client.query(oql, world_id)
return table.to_pandas(
types_mapper=pd.ArrowDtype, # 使用 Arrow-backed Pandas 类型(零拷贝)
)
# Polars 集成
def query_to_polars(client: CoomiaDipFlightClient, oql: str, world_id: str) -> pl.DataFrame:
table = client.query(oql, world_id)
return pl.from_arrow(table) # Polars 原生支持 Arrow(零拷贝)
#7. Java 服务端实现
#7.1 Quarkus Flight SQL Server
Java
@ApplicationScoped
public class CoomiaDipFlightSqlProducer implements FlightSqlProducer {
@Inject
QueryService queryService;
@Override
public FlightInfo getFlightInfoStatement(
FlightSql.CommandStatementQuery command,
CallContext context,
FlightDescriptor descriptor) {
String sql = command.getQuery();
String worldId = context.peerIdentity(); // from metadata
// 估算结果集大小
QueryPlan plan = queryService.plan(sql, worldId);
// 构建 Schema
Schema schema = toArrowSchema(plan.getOutputColumns());
// 返回 FlightInfo
return FlightInfo.builder(schema, descriptor,
List.of(new FlightEndpoint(
new Ticket(sql.getBytes(StandardCharsets.UTF_8))
)))
.setRecords(plan.getEstimatedRows())
.setBytes(plan.getEstimatedBytes())
.build();
}
@Override
public void getStreamStatement(
FlightSql.TicketStatementQuery ticket,
CallContext context,
ServerStreamListener listener) {
String sql = new String(ticket.getStatementHandle().toByteArray());
try (VectorSchemaRoot root = queryService.executeToArrow(sql)) {
listener.start(root);
// 分批发送
int batchSize = 10_000;
while (queryService.hasNext()) {
queryService.fillBatch(root, batchSize);
listener.putNext();
}
listener.completed();
} catch (Exception e) {
listener.error(e);
}
}
}
#8. 与 Doris 的集成
#8.1 Doris Arrow Flight SQL Server
Code
Doris 2.1+ 内置 Arrow Flight SQL Server:
配置启用:
arrow_flight_sql_port = 50051
enable_arrow_flight_sql = true
coomia-dip 直连 Doris Flight SQL:
Python SDK → Arrow Flight SQL → Doris FE → Doris BE (Arrow 格式)
Python
# 直连 Doris 的 Flight SQL 查询
conn = flight_sql.connect(
"grpc://doris-fe:50051",
db_kwargs={"username": "root", "password": ""},
)
cursor = conn.cursor()
cursor.execute("SELECT * FROM ontology_instances LIMIT 100000")
table = cursor.fetch_arrow_table()
# 100K 行仅需 50ms(vs JDBC 800ms)
#9. 与 Trino 的集成
#9.1 Trino Arrow Flight Connector
SQL
-- Trino 配置 Arrow Flight Connector
-- catalog/arrow.properties
connector.name=arrow-flight
arrow-flight.server.host=data-Layer
arrow-flight.server.port=50051
Python
# 通过 Trino 联邦查询,结果通过 Arrow Flight 返回
cursor.execute("""
SELECT e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d
ON e.department_id = d.id
WHERE e.world_id = 'world-001'
""")
table = cursor.fetch_arrow_table()
#10. 性能基准与优化
#10.1 传输格式对比
| 格式 | 1M 行传输时间 | CPU 开销 | 内存峰值 |
|---|---|---|---|
| JSON/REST | 12.5 s | 高 | 2.8 GB |
| JDBC ResultSet | 3.2 s | 中 | 1.5 GB |
| Arrow Flight SQL | 0.3 s | 低 | 0.4 GB |
| Arrow Flight (压缩) | 0.5 s | 低-中 | 0.3 GB |
#10.2 优化策略
Python
# 1. 批次大小调优
# 太小 → RPC 开销大;太大 → 内存峰值高
OPTIMAL_BATCH_SIZE = 65536 # 64K 行/批次
# 2. 压缩
options = flight.FlightCallOptions(
headers=[(b"x-arrow-compression", b"zstd")],
)
# 3. 投影下推(只传需要的列)
cursor.execute("SELECT id, name FROM large_table") # 不要 SELECT *
# 4. 谓词下推(在服务端过滤)
cursor.execute("SELECT * FROM table WHERE status = 'active'")
#11. Key Takeaways
| 主题 | 关键结论 |
|---|---|
| 性能 | 比 JSON/REST 快 10-100 倍 |
| 零拷贝 | Arrow 内存格式 = 传输格式 = 计算格式 |
| 柱状 | 缓存友好 + SIMD 向量化 |
| 并行 | 多 Endpoint 并行 DoGet |
| 生态 | Doris/Trino/Pandas/Polars 原生支持 |
| 流式 | RecordBatch 流式传输,内存可控 |
| 压缩 | ZSTD 压缩减少 60-80% 网络传输 |
| SDK 集成 | Python SDK 底层使用 Flight SQL |
“下一篇预告:S8-17 将深入 Pydantic v2 数据验证框架,探讨 coomia-dip 的模型层设计。