返回博客

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 的集成,以及内存管理与背压策略

#目录

  1. 为什么需要 Arrow Flight SQL
  2. Arrow 内存格式
  3. Flight RPC 协议
  4. Flight SQL 扩展
  5. coomia-dip 的数据传输架构
  6. Python 客户端实现
  7. Java 服务端实现
  8. 与 Doris 的集成
  9. 与 Trino 的集成
  10. 性能基准与优化
  11. 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/RESTArrow Flight SQL加速比
1M 行 × 10 列(数值)12.5 秒0.3 秒42x
1M 行 × 10 列(混合)18.2 秒0.8 秒23x
100K 行 × 50 列8.7 秒0.2 秒44x
流式传输 10M 行OOM4.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/REST12.5 s2.8 GB
JDBC ResultSet3.2 s1.5 GB
Arrow Flight SQL0.3 s0.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 的模型层设计。