InfluxElixir.Flight.Reader (InfluxElixir v0.1.35)

Copy Markdown View Source

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 IPC Message flatbuffer. The first message in a stream contains a Schema message; subsequent messages contain RecordBatch messages 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 TypeElixir value
Int8-64, UInt8-64integer()
Float16/32/64float()
Boolboolean()
Utf8, LargeUtf8, Utf8View (string functions return it)binary()
Binary, LargeBinary, BinaryViewlowercase hex, as HTTP renders it
TimestampDateTime.t() (microsecond precision, as on HTTP)
Date32/64"YYYY-MM-DD"
Duration"PT60S", "PT0.5S", "-PT0.000000001S", "P0D"
Decimal128/256a number (integer at scale 0, else float)
Struct (selector_* without a subscript)a map of its members
List, LargeList, FixedSizeList (array_agg)a list
Nullnil

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

column_schema()

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

kind()

@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

decode_flight_data(list)

@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 of FlightData structs from a DoGet stream

Example

iex> InfluxElixir.Flight.Reader.decode_flight_data([])
{:ok, []}

parse_schema(header)

@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}.