Arrow IPC record batch decoder for Arrow Flight query results.
Converts a list of FlightData messages (received from a DoGet gRPC
stream) into a list of Elixir row maps.
Arrow IPC Format Overview
Each FlightData message carries two binary blobs:
data_header— serialised Arrow IPCMessageflatbuffer. The first message in a stream contains aSchemamessage; subsequent messages containRecordBatchmessages with buffer offset/length metadata.data_body— raw column buffer bytes referenced by the batch metadata.
Schema and record batch metadata is parsed using a proper FlatBuffer
binary reader (InfluxElixir.Flight.FlatBuffer), following the Arrow
IPC FlatBuffer schema specification exactly.
Supported Column Types
Every type decodes to the value the HTTP transport returns for it, so a
row is the same map on both (verified against InfluxDB 3, and recorded in
test/fixtures/flight):
| Arrow Type | Elixir value |
|---|---|
| Int8-64, UInt8-64 | integer() |
| Float16/32/64 | float() |
| Bool | boolean() |
| Utf8, LargeUtf8, Utf8View (string functions return it) | binary() |
| Binary, LargeBinary, BinaryView | lowercase hex, as HTTP renders it |
| Timestamp | DateTime.t() (microsecond precision, as on HTTP) |
| Date32/64 | "YYYY-MM-DD" |
| Duration | "PT60S", "PT0.5S", "-PT0.000000001S", "P0D" |
| Decimal128/256 | a number (integer at scale 0, else float) |
Struct (selector_* without a subscript) | a map of its members |
List, LargeList, FixedSizeList (array_agg) | a list |
| Null | nil |
Nested types are read with the batch's field nodes and, for view types, its variadic buffer counts, in Arrow's depth-first layout.
Null bitmaps are supported. A null cell is left out of the row map, the
same as a null column in InfluxDB 3's JSON, so a row is identical over
Flight and HTTP; assert with refute Map.has_key?(row, "col").
Limitations
Interval, Time, Map, Union, FixedSizeBinary, run-end encoded and
list-view columns, dictionary-encoded fields and compressed IPC bodies
are refused with {:error, {:unsupported_arrow_type, type, column}}
rather than dropped. (InfluxDB 3 sends tag columns hydrated, so their
Dictionary(Int32, Utf8) type does not reach the reader.)
Summary
Types
Parsed column schema entry (unit is set for Timestamp columns)
How a column decodes (see "Supported Column Types").
Functions
Decodes a list of FlightData messages into row maps.
Extracts column name/type pairs from an Arrow IPC Schema message header.
Types
@type column_schema() :: %{ name: binary(), type_id: non_neg_integer(), unit: System.time_unit() | nil, kind: kind(), children: [column_schema()] }
Parsed column schema entry (unit is set for Timestamp columns)
@type kind() :: :primitive | :null | :binary | :large_binary | :large_utf8 | :binary_view | :utf8_view | :float16 | :struct | :list | :large_list | {:fixed_size_list, integer()} | {:duration, System.time_unit()} | {:date, :day | :millisecond} | {:decimal, integer(), integer()} | {:unsupported, binary()}
How a column decodes (see "Supported Column Types").
Functions
@spec decode_flight_data([InfluxElixir.Flight.Proto.FlightData.t()]) :: {:ok, [map()]} | {:error, term()}
Decodes a list of FlightData messages into row maps.
The first element of flight_data_list is expected to be the schema message
(typically with an empty data_body). Subsequent elements are record batch
messages.
Returns {:ok, [map()]} on success or {:error, reason} on parse failure.
Parameters
flight_data_list— ordered list ofFlightDatastructs from a DoGet stream
Example
iex> InfluxElixir.Flight.Reader.decode_flight_data([])
{:ok, []}
@spec parse_schema(binary() | nil) :: {:ok, [column_schema()]} | {:error, term()}
Extracts column name/type pairs from an Arrow IPC Schema message header.
Parses the FlatBuffer metadata according to the Arrow IPC specification: Message → Schema → Field[] → name + Type union.
Returns {:ok, [column_schema()]} or {:error, reason}.