Skip to content
Roboto
Esc
↑↓navigate↵open⌘Jpreview
On this page

roboto.experimental.topics.decode.common

What reading one file of a read plan’s partition takes and gives.

Module Contents

FileDecodeParams

class roboto.experimental.topics.decode.common.FileDecodeParams#View Source

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.

Attributes

FileDecodeParams.cache_dir

cache_dir pathlib.Path #

Directory Parquet files are cached under.

FileDecodeParams.cache_policy

Whether fetched Parquet files are cached to local disk.

FileDecodeParams.signed_url_resolver

signed_url_resolver SignedUrlResolver #

Mints a signed download URL for a file.

FileDecoder

class roboto.experimental.topics.decode.common.FileDecoder#View Source

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.

FileDecoder.batches()

abstract batches()#View Source

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

collections.abc.Iterator[pyarrow.RecordBatch]

FileDecoder.close()

abstract close()#View Source

Release the file. Safe to call more than once.

Return type

None

FileDecoder.struct_field_names()

abstract struct_field_names(path)#View Source

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.

Return type

Optional[list[str]]

Properties

FileDecoder.value_fields

value_fields list[pyarrow.Field] abstract #

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]

FileDecoderOpener

roboto.experimental.topics.decode.common.FileDecoderOpener#View Source

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

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

ScanTaskGroup

class roboto.experimental.topics.decode.common.ScanTaskGroup#View Source

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.

Attributes

ScanTaskGroup.supplies

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.

ScanTaskGroup.topic_name

topic_name str | None #

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

SignedUrlResolver

roboto.experimental.topics.decode.common.SignedUrlResolver#View Source

Resolves a file id (fs_node_id) to a signed download URL.

SuppliedField

class roboto.experimental.topics.decode.common.SuppliedField#View Source

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

Attributes

SuppliedField.excluded

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

Paths strictly inside path, none inside another.

Was this page helpful?