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_KEYandROW_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
_indexindex thatTopic.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
batchwithout 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.xalongside a structposewith childx. A plain dict would silently drop one (last write wins); the ambiguity is rejected instead. Rename the offending field or disableflatten=Trueto 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_KEYmetadata isb"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_KEYmetadata isb"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