roboto.experimental.topics.decode#

Decoding the files of a read plan’s partitions into RecordBatches of topic data.

Submodules#

Package Contents#

roboto.experimental.topics.decode.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.FileDecodeParams#

What decoding a file takes beyond the read plan: how to reach the file and whether to cache it.

Caching applies to Parquet files only; MCAP files always stream.

cache_dir: pathlib.Path#

Directory Parquet files are cached under.

cache_policy: roboto.storage.CachePolicy#

Whether fetched Parquet files are cached to local disk.

signed_url_resolver: SignedUrlResolver#

Mints a signed download URL for a file.

class roboto.experimental.topics.decode.FileDecoder#

Bases: abc.ABC

Decodes the fields one file supplies to a partition into RecordBatches.

Opening a decoder opens its file, so its fields are known before the first batch. Close it when done, or use it as a context manager.

abstract 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]

abstract close()#

Release the file. Safe to call more than once.

Return type:

None

abstract 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]#
Abstractmethod:

Return type:

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.

roboto.experimental.topics.decode.FileDecoderOpener#

Opens a FileDecoder of a group’s file for its partition and the plan’s window.

The window is absolute and includes both ends. A decoder keeps the rows whose absolute timestamp lies in it.

class roboto.experimental.topics.decode.ScanTaskGroup#

The scan tasks of a partition that read one file, and the fields they supply.

Every scan task in the group agrees on the file’s format, transformations and topic name.

format: roboto.domain.topics.RepresentationStorageFormat#
object: roboto.experimental.topics.read_plan.ReadPlanObjectRef#
supplies: tuple[SuppliedField, ...]#

In the order of the leaf-most paths they come from.

For one path, the supplied field at the path itself comes before the subtrees inside it, which come in plan order. Empty when the file is read only for its row numbers and timestamps.

topic_name: str | None#

The topic the scan tasks read from the file; see topic_name.

class roboto.experimental.topics.decode.SuppliedField#

The field at path, which one file supplies to the read, less the fields at excluded, which other scan tasks supply (they may read the same file).

excluded: tuple[roboto.domain.topics.record.FieldPath, ...] = ()#

Paths strictly inside path, none inside another.

path: roboto.domain.topics.record.FieldPath#
roboto.experimental.topics.decode.make_file_decoder_opener(params)#

Return a FileDecoderOpener that reads files through params.

The opener picks the decoder by the group’s format: McapFileDecoder or ParquetFileDecoder.

Parameters:

params (roboto.experimental.topics.decode.common.FileDecodeParams) – How to reach each file and whether to cache it.

Returns:

An opener of file decoders. The opener raises RobotoReadPlanExecutionException with kind unsupported-format for a format other than MCAP or Parquet, and passes on what opening the decoder raises.

Return type:

roboto.experimental.topics.decode.common.FileDecoderOpener