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

roboto.experimental.topics.batch_transforms

Representation conversion for topic-data RecordBatches.

Topic data moves through the read path as Arrow RecordBatches in its public shape: one column per top-level projected field, with struct/list types mirroring the schema tree, plus one dedicated timestamp column of absolute Unix-epoch nanoseconds (int64) marked by field metadata (TIMESTAMP_FIELD_METADATA_KEY).

Inside the read path, a decoded batch also leads with a row number column (topic_data_schema()): each row’s uint64 position among its file’s rows of the topic, marked by ROW_NUMBER_FIELD_METADATA_KEY. When a partition’s fields are stored across several files, merging those files compares this column to confirm every file holds the same rows. Batches returned to a caller leave it out (drop_row_number_column()).

flatten_table() expands a table’s struct columns into dot-delimited leaf columns, which Topic.get_data_as_df(flatten=True) returns as the value columns of its DataFrame.

It also exposes the helpers that construct and locate the timestamp and row number columns (topic_data_schema(), timestamp_field(), timestamp_column_index(), row_number_column_index()), which the read path uses to mark and find those columns by metadata rather than name.

Module Contents

MARKER_FIELD_METADATA_VALUE

roboto.experimental.topics.batch_transforms.MARKER_FIELD_METADATA_VALUE = b'true'#View Source

Value of the metadata key on the column it marks, for both TIMESTAMP_FIELD_METADATA_KEY and ROW_NUMBER_FIELD_METADATA_KEY.

ROW_NUMBER_FIELD_METADATA_KEY

roboto.experimental.topics.batch_transforms.ROW_NUMBER_FIELD_METADATA_KEY = b'roboto.topic_data.row_number'#View Source

Arrow field-metadata key marking the row number column of a decoded topic-data batch.

The row number column holds this key with the value b"true" (MARKER_FIELD_METADATA_VALUE); a column holding the key with any other value is a value column.

ROW_NUMBER_FIELD_NAME

roboto.experimental.topics.batch_transforms.ROW_NUMBER_FIELD_NAME = '_row'#View Source

Requested name of the row number column; _ is appended while another column of the batch has it.

The column’s identity is its metadata marker (ROW_NUMBER_FIELD_METADATA_KEY), never this name.

TIMESTAMP_FIELD_METADATA_KEY

roboto.experimental.topics.batch_transforms.TIMESTAMP_FIELD_METADATA_KEY = b'roboto.topic_data.timestamp'#View Source

Arrow field-metadata key marking the per-row timestamp column of a topic-data batch.

The timestamp column holds this key with the value b"true" (MARKER_FIELD_METADATA_VALUE); a column holding the key with any other value is a value column.

TIMESTAMP_FIELD_NAME

roboto.experimental.topics.batch_transforms.TIMESTAMP_FIELD_NAME = '_index'#View Source

Name of the emitted per-row timestamp column.

Source-neutral by design: the column always carries the resolved timeline’s absolute Unix-epoch nanoseconds, whatever that source is (message log time, publish time, or a schema field), so the name asserts no particular origin. It matches the _index index that Topic.get_data_as_df() labels its rows with. The column’s real identity is its metadata marker (TIMESTAMP_FIELD_METADATA_KEY), never this name, which is uniquified by suffixing when a projected field already claims it.

drop_row_number_column()

roboto.experimental.topics.batch_transforms.drop_row_number_column(batch)#View Source

Return batch without its row number column, found by its metadata marker.

Parameters

batch pyarrow.RecordBatch

Raises

The batch does not contain exactly one marked row number column.

Return type

pyarrow.RecordBatch

flatten_table()

roboto.experimental.topics.batch_transforms.flatten_table(table)#View Source

Expand struct columns into dot-delimited leaf columns, recursively.

A null at any struct level propagates to nulls in every leaf column beneath it. List-typed columns stay whole. This is the DataFrame packing shape: dotted leaf columns over the projected tree.

Parameters

table pyarrow.Table

Raises

Two columns resolve to the same dotted name — e.g. a top-level field literally named pose.x alongside a struct pose with child x. A plain dict would silently drop one (last write wins); the ambiguity is rejected instead. Rename the offending field or disable flatten=True to recover the column.

Return type

pyarrow.Table

row_number_column_index()

roboto.experimental.topics.batch_transforms.row_number_column_index(schema)#View Source

Locate the row number column: the field whose ROW_NUMBER_FIELD_METADATA_KEY metadata is b"true".

Parameters

schema pyarrow.Schema

Raises

The schema does not contain exactly one marked column.

Return type

int

timestamp_column_index()

roboto.experimental.topics.batch_transforms.timestamp_column_index(schema)#View Source

Locate the timestamp column: the field whose TIMESTAMP_FIELD_METADATA_KEY metadata is b"true".

The column is identified by metadata, never by name: a projected root field can legitimately carry any name, including the timestamp column’s conventional one.

Parameters

schema pyarrow.Schema

Raises

The schema does not contain exactly one marked column.

Return type

int

timestamp_field()

roboto.experimental.topics.batch_transforms.timestamp_field(name=TIMESTAMP_FIELD_NAME)#View Source

The timestamp column’s Arrow field: int64 epoch nanoseconds, metadata-marked.

Parameters

name str

Return type

pyarrow.Field

topic_data_schema()

roboto.experimental.topics.batch_transforms.topic_data_schema(value_fields)#View Source

The schema of a decoded batch whose value columns are value_fields: row number, timestamp, then values.

The row number column is non-null uint64. _ is appended to the timestamp column’s name until no value column has it, then to the row number column’s name until neither a value column nor the timestamp column has it.

Parameters

value_fields collections.abc.Sequence[pyarrow.Field]

Return type

pyarrow.Schema

Was this page helpful?