Back to Blog

Arrow Flight SQL Deep Dive: High-Performance Columnar Data Transport

1. [Why Arrow Flight SQL](#1-why-arrow-flight-sql)

CoomiaPublished on November 25, 20257 min read
Share this articleTwitter / X

Series: S8 Technology Deep Dives · Article 16 | Level: Advanced | Reading Time: 20 min

Arrow Flight SQL Deep Dive: High-Performance Columnar Data Transport

#TL;DR

  • Apache Arrow Flight SQL is the key technology in coomia-dip for high-performance data transport, passing query results between services in zero-copy columnar format, 10-100x faster than JSON serialization
  • This article deeply analyzes the Arrow memory format, Flight RPC protocol, Flight SQL extensions, and coomia-dip's high-speed data channel from Data Layer to SDK
  • Covers Arrow IPC format, RecordBatch streaming, DoPut/DoGet protocol, Doris/Trino integration, and memory management with backpressure strategies

#Table of Contents

  1. Why Arrow Flight SQL
  2. Arrow Memory Format
  3. Flight RPC Protocol
  4. Flight SQL Extension
  5. coomia-dip Data Transport Architecture
  6. Python Client Implementation
  7. Java Server Implementation
  8. Doris Integration
  9. Trino Integration
  10. Performance Benchmarks and Optimization
  11. Key Takeaways

#1. Why Arrow Flight SQL

#1.1 Traditional Data Transport Bottlenecks

Code
Traditional JSON/REST transport:
  Client ← JSON ← Server
    ├── Serialization: Object → JSON String (CPU intensive)
    ├── Network: JSON text 3-5x larger than binary
    ├── Deserialization: JSON String → Object (CPU intensive)
    └── Memory: One object per row → GC pressure

Arrow Flight SQL transport:
  Client ← Arrow IPC ← Server
    ├── Zero serialization: Memory format = wire format
    ├── Network efficient: Binary columnar, compact
    ├── Zero copy: Direct memory mapping
    └── SIMD friendly: Columnar layout for vectorized compute

#1.2 Performance Comparison

ScenarioJSON/RESTArrow Flight SQLSpeedup
1M rows x 10 cols (numeric)12.5s0.3s42x
1M rows x 10 cols (mixed)18.2s0.8s23x
100K rows x 50 cols8.7s0.2s44x
Stream 10M rowsOOM4.2s-

#2. Arrow Memory Format

#2.1 Columnar Layout

Code
Row-oriented (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}

Columnar (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

  Each column stored contiguously → Cache friendly → SIMD vectorization

#2.2 Buffer Structure

Code
Int64 Column "id" (3 values):
┌──────────────┬────────────────────────────┐
│ Validity Bitmap │ Data Buffer (24 bytes)  │
│ [1, 1, 1]       │ [01 00 00 00 00 00 00 00│
│                  │  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 and Schema

Python
import pyarrow as pa

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()),
])

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,
)

#3. Flight RPC Protocol

#3.1 Protocol Overview

Code
Flight protocol is built on gRPC with 4 core RPCs:

1. GetFlightInfo(FlightDescriptor) → FlightInfo
   "How big is this dataset? Where is it?"

2. DoGet(Ticket) → Stream<FlightData>
   "Give me the data" (server streaming)

3. DoPut(Stream<FlightData>) → PutResult
   "Accept data" (client streaming)

4. DoExchange(Stream<FlightData>) → Stream<FlightData>
   "Bidirectional data exchange" (bidi streaming)

#3.2 Data Flow

Code
Client                          Flight Server
  │── GetFlightInfo(descriptor) ────→│
  │←── FlightInfo(endpoints, size) ──│
  │── DoGet(ticket) ────────────────→│
  │←── Schema ───────────────────────│
  │←── RecordBatch 1 ───────────────│
  │←── RecordBatch N ───────────────│
  │←── (stream end) ────────────────│

#3.3 Parallel Transport

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 receives in parallel

#4. Flight SQL Extension

Code
Flight SQL = Flight + SQL semantics

Key RPC extensions:
  GetFlightInfo + CommandStatementQuery  → Execute SQL query
  DoGet                                  → Fetch query results
  DoPut + CommandStatementUpdate        → Execute DML
  GetCatalogs / GetSchemas / GetTables  → Metadata queries

#5. coomia-dip Data Transport Architecture

Code
┌────────────────────────────────────────┐
│              Python SDK                 │
│  ADBC Flight SQL Client                │
│  → Arrow RecordBatch → Pandas / Polars │
└────────────────┬───────────────────────┘
                 │ Flight SQL (gRPC)
┌────────────────┼───────────────────────┐
│  Data Layer    │                       │
│  Flight SQL Server (Quarkus)           │
│  → SQL parsing → Route to storage      │
│       ┌────────┴────────┐              │
│   Doris (OLAP)    Trino (Federation)   │
└────────────────────────────────────────┘

#6. Python Client Implementation

Python
import pyarrow.flight as flight

class CoomiaDipFlightClient:
    def __init__(self, host: str, port: int = 50051):
        self.client = flight.FlightClient(f"grpc://{host}:{port}")

    def query(self, sql: str, world_id: str) -> pa.Table:
        descriptor = flight.FlightDescriptor.for_command(sql.encode())
        flight_info = self.client.get_flight_info(
            descriptor,
            options=flight.FlightCallOptions(
                headers=[(b"x-world-id", world_id.encode())],
            ),
        )
        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):
        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

# Pandas integration (zero-copy)
def query_to_pandas(client, oql, world_id) -> pd.DataFrame:
    table = client.query(oql, world_id)
    return table.to_pandas(types_mapper=pd.ArrowDtype)

# Polars integration (zero-copy)
def query_to_polars(client, oql, world_id) -> pl.DataFrame:
    return pl.from_arrow(client.query(oql, world_id))

#7. Java Server Implementation

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();
        QueryPlan plan = queryService.plan(sql);
        Schema schema = toArrowSchema(plan.getOutputColumns());
        return FlightInfo.builder(schema, descriptor,
            List.of(new FlightEndpoint(new Ticket(sql.getBytes()))))
            .setRecords(plan.getEstimatedRows())
            .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);
            while (queryService.hasNext()) {
                queryService.fillBatch(root, 65536);
                listener.putNext();
            }
            listener.completed();
        }
    }
}

#8. Doris Integration

Code
Doris 2.1+ built-in Arrow Flight SQL Server:
  arrow_flight_sql_port = 50051
  enable_arrow_flight_sql = true

Direct connection: Python SDK → Arrow Flight SQL → Doris
100K rows in 50ms (vs JDBC 800ms)

#9. Trino Integration

SQL
-- Trino Arrow Flight Connector
-- catalog/arrow.properties
connector.name=arrow-flight
arrow-flight.server.host=data-Layer
arrow-flight.server.port=50051

#10. Performance Benchmarks and Optimization

Format1M Row TransferCPU OverheadPeak Memory
JSON/REST12.5sHigh2.8 GB
JDBC ResultSet3.2sMedium1.5 GB
Arrow Flight SQL0.3sLow0.4 GB
Arrow Flight (compressed)0.5sLow-Med0.3 GB

Optimization strategies: Batch size tuning (64K rows optimal), ZSTD compression, projection pushdown, predicate pushdown.

#11. Key Takeaways

TopicKey Conclusion
Performance10-100x faster than JSON/REST
Zero copyArrow memory format = wire format = compute format
ColumnarCache friendly + SIMD vectorization
ParallelMulti-endpoint parallel DoGet
EcosystemDoris/Trino/Pandas/Polars native support
StreamingRecordBatch streaming with controlled memory
CompressionZSTD reduces network transfer by 60-80%
SDK integrationPython SDK uses Flight SQL underneath

Next up: S8-17 dives into Pydantic v2 data validation framework, exploring coomia-dip's model layer design.