roboto.experimental.topics.decode.parquet
Module Contents
CACHED_PARQUET_NAME_PATTERN
Filename template for locally cached Parquet files; keyed on the stable file id.
ParquetFileDecoder
Bases: roboto.experimental.topics.decode.common.FileDecoder
Decodes one Parquet file of a partition, reading only the columns that hold its supplied fields and timestamp.
A field the group supplies with fields excluded is expanded into its other fields when it is a struct reached through structs only, and so is each struct inside it that holds an excluded field. Any other field that holds an excluded field, such as a list or a map, is kept whole, excluded field included. A struct left with no field is dropped.
A TIMESTAMP field counts in its own unit, and any other numeric field in the plan’s unit, nanoseconds when the plan gives none. Row groups whose timestamp statistics, shifted by the partition’s time_offset_ns, rule the window out are not read. Every timestamp of a row group that is read must have a signed 64-bit nanosecond value once shifted, whether or not the window keeps its row.
A partition that declares a data_range owns only that half-open span of the file’s stored rows, counted from 0. Position in the file, not time, separates the partition’s rows from the file’s other slices, so the slice applies before the window. Row groups entirely outside the slice are not read, a row group that straddles a boundary is read and trimmed, and only the timestamps of rows inside the slice must have a signed 64-bit nanosecond value. See data_range for the slice contract.
ParquetFileDecoder.batches()
The window’s rows, in the file’s stored row order; iterate it once.
A partition that declares a data_range gets only the window’s rows inside that slice of the file. Each batch has the columns of topic_data_schema() over value_fields: the row number, the timestamp, then the value columns. A row’s number is its 0-based position among the file’s rows of the topic, counting every stored row, including rows outside the window or the data_range and rows with a null timestamp, so a row has the same number in every file of its partition. The timestamp is absolute: the stored value in nanoseconds plus the partition’s time_offset_ns. Batch boundaries carry no meaning.
Return type
ParquetFileDecoder.close()
Release the file. Safe to call more than once.
Return type
ParquetFileDecoder.struct_field_names()
The names of the fields of the struct at path in the file, in the file’s order.
None when the file has no struct at path.
Parameters
Return type
Properties
ParquetFileDecoder.value_fields
The value columns, one per top-level field the file supplies, sorted by name, comparing Unicode code points.
Each struct keeps the fields the file supplies, in the order the file stores them.