roboto.experimental.topics.split_partition#

Combine the files of a split partition into the partition’s rows.

Each combined row holds every projected top-level field, and each file gives the fields it supplies.

Module Contents#

class roboto.experimental.topics.split_partition.SplitPartition(topic_part_id, file_decoders, top_level_names)#

A split partition, read through one file decoder per scan task group.

A split partition’s scan tasks, grouped by file, form more than one scan task group, so its rows are read from more than one file. Several scan tasks on one file are one group, read as one file. The groups, and so the file decoders, are ordered by the lowest precedence among their scan tasks.

Each top-level field is made from the files’ decoded columns by these rules:

  • A field one file decoded is taken whole from it.

  • A field several files decoded must be a struct in each of them, and is a struct combined from theirs. Its fields are the names struct_field_names() gives for the struct in those files, in the order the first of them stores them, then any name that file lacks, in the next such file’s order. A name none of those files decoded is dropped, and each other field is made by these same rules from the files that decoded it. The combined struct is valid in a row where any of those files has it valid.

  • A top-level name no file decoded is an empty struct, null in every row.

Parameters:
combine(file_batches)#

Combine each file decoder’s rows, paired in stored order, into the partition’s rows.

Row numbers and timestamps come from the file of the first scan task group. The combined batches have the columns of topic_data_schema() over value_fields.

Parameters:

file_batches (collections.abc.Sequence[collections.abc.Sequence[pyarrow.RecordBatch]]) – Each file decoder’s batches, in the order of the file decoders.

Raises:

RobotoReadPlanExecutionException – With kind scan-task-row-mismatch, when the files do not hold the same rows: the same row numbers and timestamps in the same order. The message names the first row that differs, including when one file runs out of rows before another.

Return type:

collections.abc.Iterator[pyarrow.RecordBatch]

property value_fields: list[pyarrow.Field]#

The combined partition’s value columns, one per top-level name, in the order given.

Return type:

list[pyarrow.Field]