Arrow Flight
Everything so far has been about avoiding copies within one process (Lessons 1–3) or between memory and disk (Lesson 7). Flight is the same idea applied across a network: RecordBatches go out as Arrow IPC-formatted streams over gRPC, with no per-row (de)serialization.
- Flight is a gRPC-based streaming protocol, not a request/response REST API. Client and server exchange streams of RecordBatches.
- A
FlightDescriptoridentifies a dataset (by path or an opaque command);FlightInfodescribes it (schema, row count, where to fetch it — aTicket); aTicketis what you hand todo_getto actually retrieve the data. do_get(server → client stream) anddo_put(client → server stream) are the two core data-movement RPCs;list_flights/get_flight_infoare metadata/discovery RPCs.- In Python: implement a server by subclassing
pyarrow.flight.FlightServerBaseand overriding the RPCs you support; talk to it withpyarrow.flight.FlightClient.
Exercise
import threading
import time
import pyarrow as pa
import pyarrow.flight as flight
table = pa.table({
"id": [1, 2, 3],
"amount": [100, 150, 200],
})
class DemoFlightServer(flight.FlightServerBase):
def __init__(self, location, table):
super().__init__(location)
self._table = table
def do_get(self, context, ticket):
return flight.RecordBatchStream(self._table)
def list_flights(self, context, criteria):
descriptor = flight.FlightDescriptor.for_path(b"demo")
info = flight.FlightInfo(self._table.schema, descriptor, [], self._table.num_rows, 0)
yield info
location = "grpc://0.0.0.0:0"
server = DemoFlightServer(location, table)
server_thread = threading.Thread(target=server.serve, daemon=True)
server_thread.start()
time.sleep(0.3)
client = flight.FlightClient(f"grpc://127.0.0.1:{server.port}")
reader = client.do_get(flight.Ticket(b"demo"))
result_table = reader.read_all()
print("Client received rows:", result_table.num_rows)
print("Round-trip equal to original:", result_table.equals(table))
server.shutdown()
Client received rows: 3
Round-trip equal to original: True
do_get response before doing anything with it defeats the point of a streaming protocol; large results are meant to be consumed batch-by-batch as they arrive. Also: forgetting to override the RPC you actually need (e.g. calling do_get on a server that never overrode it) gets you an unhelpful "not implemented" failure instead of a clear error about a missing method.
Retrieval check
What does a client actually hand to do_get to retrieve data from a Flight server?
Correct. Ticket is what do_get consumes to actually retrieve the data stream — FlightDescriptor/FlightInfo are how you discover which ticket to use.
Not quite. do_get takes a Ticket. The descriptor/info pair is how a client discovers what ticket to ask for.
Is Arrow Flight a request/response REST-style API?
Correct. Flight is fundamentally a stream of RecordBatches over gRPC, not a REST request/response model.
No. Flight is a gRPC-based streaming protocol — treating it like REST (e.g. buffering the whole response) defeats its purpose.
Practice
- Extend the server to serve a Table you built earlier in this course (Lesson 5 or 7's table works well).
- Fetch it with a client on a different
Ticketvalue. - Confirm
.equals(...)returnsTrue. - Explain in one sentence what would be serialized instead of shared if this were plain gRPC without Arrow (i.e., row-by-row protobuf messages).
New terms — Flight — are in the glossary.
do_put look like for this same server, accepting a Table from a client instead of serving one?" Next lesson sees what Flight looks like specialized for SQL workloads: Arrow Flight SQL.