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#

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

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

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

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.

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

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.

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

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.

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

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.

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

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

Raises:

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

Parameters:

batch (pyarrow.RecordBatch)

Return type:

pyarrow.RecordBatch

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

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.

Raises:

RobotoInvalidRequestException – 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.

Parameters:

table (pyarrow.Table)

Return type:

pyarrow.Table

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

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

Raises:

RobotoInternalException – The schema does not contain exactly one marked column.

Parameters:

schema (pyarrow.Schema)

Return type:

int

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

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.

Raises:

RobotoInternalException – The schema does not contain exactly one marked column.

Parameters:

schema (pyarrow.Schema)

Return type:

int

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

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

Parameters:

name (str)

Return type:

pyarrow.Field

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

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