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
- Why Arrow Flight SQL
- Arrow Memory Format
- Flight RPC Protocol
- Flight SQL Extension
- coomia-dip Data Transport Architecture
- Python Client Implementation
- Java Server Implementation
- Doris Integration
- Trino Integration
- Performance Benchmarks and Optimization
- 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
| Scenario | JSON/REST | Arrow Flight SQL | Speedup |
|---|---|---|---|
| 1M rows x 10 cols (numeric) | 12.5s | 0.3s | 42x |
| 1M rows x 10 cols (mixed) | 18.2s | 0.8s | 23x |
| 100K rows x 50 cols | 8.7s | 0.2s | 44x |
| Stream 10M rows | OOM | 4.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
| Format | 1M Row Transfer | CPU Overhead | Peak Memory |
|---|---|---|---|
| JSON/REST | 12.5s | High | 2.8 GB |
| JDBC ResultSet | 3.2s | Medium | 1.5 GB |
| Arrow Flight SQL | 0.3s | Low | 0.4 GB |
| Arrow Flight (compressed) | 0.5s | Low-Med | 0.3 GB |
Optimization strategies: Batch size tuning (64K rows optimal), ZSTD compression, projection pushdown, predicate pushdown.
#11. Key Takeaways
| Topic | Key Conclusion |
|---|---|
| Performance | 10-100x faster than JSON/REST |
| Zero copy | Arrow memory format = wire format = compute format |
| Columnar | Cache friendly + SIMD vectorization |
| Parallel | Multi-endpoint parallel DoGet |
| Ecosystem | Doris/Trino/Pandas/Polars native support |
| Streaming | RecordBatch streaming with controlled memory |
| Compression | ZSTD reduces network transfer by 60-80% |
| SDK integration | Python SDK uses Flight SQL underneath |
“Next up: S8-17 dives into Pydantic v2 data validation framework, exploring coomia-dip's model layer design.