roboto.experimental.topics.decode.parquet#

Module Contents#

roboto.experimental.topics.decode.parquet.CACHED_PARQUET_NAME_PATTERN = '{fs_node_id}.parquet'#

Filename template for locally cached Parquet files; keyed on the stable file id.

class roboto.experimental.topics.decode.parquet.ParquetFileDecoder(group, partition, window, params)#

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.

Parameters:
batches()#

The window’s rows, in the file’s stored row order; iterate it once.

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 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:

collections.abc.Iterator[pyarrow.RecordBatch]

close()#

Release the file. Safe to call more than once.

Return type:

None

struct_field_names(path)#

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:

path (roboto.domain.topics.record.FieldPath)

Return type:

Optional[list[str]]

property value_fields: list[pyarrow.Field]#

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.

Return type:

list[pyarrow.Field]