roboto.experimental.topics.plan_execution#
Execute a read plan: decode each partition’s files and yield the topic’s rows as RecordBatches.
execute_read_plan() lists the steps of a read. This module holds the projected paths, the leaf-most paths,
assigning scan tasks and grouping them by file, and reading the partitions in plan order.
Decoding each file lives in decode,
combining a split partition in split_partition,
and the output schema in that module and batch_transforms.
Module Contents#
- roboto.experimental.topics.plan_execution.SchemaFieldPaths#
Returns the path of every field a schema declares, given the schema’s id.
- roboto.experimental.topics.plan_execution.assign_scan_tasks(partition, leaf_most)#
Return the groups of
partition’s scan tasks that supply the leaf-most pathsleaf_most.A path’s supplier is the highest-precedence scan task whose subtree contains it; a task without a subtree contains every field. Ties go to the later task in plan order. A task whose subtree lies strictly inside a path, and which is the supplier of its own subtree, supplies that subtree. Each supplied field excludes only the outermost supplied fields strictly inside it.
The tasks on one file form one group whatever their precedence, so each file is opened once. Groups come in order of their lowest precedence, ties in the order each file first appears among the scan tasks. A group that supplies nothing is dropped, unless every group supplies nothing; then the first group is kept, and its file gives only the row number and timestamp columns.
- Raises:
RobotoReadPlanExecutionException – With kind
inconsistent-scan-tasks-on-file, when scan tasks on one file disagree on its format, transformations or topic name; with kindprojected-field-in-no-scan-task, when no scan task contains a path ofleaf_most.- Parameters:
partition (roboto.experimental.topics.read_plan.ReadPlanPartition)
leaf_most (collections.abc.Sequence[roboto.domain.topics.record.FieldPath])
- Return type:
list[roboto.experimental.topics.decode.common.ScanTaskGroup]
- roboto.experimental.topics.plan_execution.execute_read_plan(plan, projected, open_file_decoder)#
Decode the files a read plan names and yield the topic’s rows as RecordBatches.
A read takes these steps:
Projected paths:
projected, fromprojected_paths().Leaf-most paths (
leaf_most_paths()).Assign scan tasks (
assign_scan_tasks()).Group by file (
assign_scan_tasks()).Decode each file, through the decoders
open_file_decoderopens.Combine a split partition (
SplitPartition).Output schema (
topic_data_schema()).Partitions: read in plan order, each one’s schema checked against the first partition’s when its files open, before any of its rows.
Partitions without scan tasks are skipped. Partitions are yielded in plan order and their rows are never interleaved; within a partition, rows keep their stored order. Nothing is sorted by time or deduplicated, so a consumer that needs rows in time order sorts them.
A plan with one partition to read yields its rows as they are decoded. A plan with several decodes up to 32 partitions at once and holds each partition’s rows until it is yielded. A split partition’s files are decoded in full before their rows are combined.
- Parameters:
plan (roboto.experimental.topics.read_plan.ReadPlan) – The read plan the service resolved.
projected (collections.abc.Sequence[roboto.domain.topics.record.FieldPath]) – The field paths the plan projects, from
projected_paths().open_file_decoder (roboto.experimental.topics.decode.common.FileDecoderOpener) – Opens the decoder of one scan task group’s file.
- Yields:
RecordBatches with the columns of
topic_data_schema()– the row number, the timestamp in absolute Unix-epoch nanoseconds, then one value column per projected top-level field, sorted by name comparing Unicode code points. A batch holds the rows of one partition only; batch sizes and boundaries are otherwise arbitrary.- Raises:
RobotoReadPlanExecutionException – With the kind of the step that refuses the plan:
projected-field-in-no-scan-taskorinconsistent-scan-tasks-on-file(steps 3 and 4);field-split-inside-non-structorscan-task-row-mismatch(step 6);partition-schema-mismatch, when a partition’s files give the read a different schema than the first partition’s, raised before any error in that partition’s rows and even when it has no rows in the window (step 8); and the kindsopen_file_decoderand its decoders raise, such asunsupported-format,unsupported-timestamp,invalid-timestampandfield-not-in-file(step 5).- Return type:
collections.abc.Generator[pyarrow.RecordBatch, None, None]
- roboto.experimental.topics.plan_execution.leaf_most_paths(paths)#
Return the paths among
pathsthat have no descendant among them.A projection lists a struct and its children as separate fields; reading the struct would read every child, including one the projection leaves out, so only the leaf-most paths are read. The empty path, which names the schema root, and duplicates are dropped. The paths are sorted by their components, each compared by Unicode code point.
- Parameters:
paths (collections.abc.Iterable[roboto.domain.topics.record.FieldPath])
- Return type:
list[roboto.domain.topics.record.FieldPath]
- roboto.experimental.topics.plan_execution.projected_paths(plan, schema_field_paths)#
Return the field paths
planprojects.A projection that lists its fields gives their paths. A projection of every field (
projection.all) gives every field the plan’s schema declares, fetched throughschema_field_paths. This is the only case that fetches them, and a plan with no scan task in any partition fetches nothing and gives no paths.- Parameters:
plan (roboto.experimental.topics.read_plan.ReadPlan) – The read plan the service resolved.
schema_field_paths (SchemaFieldPaths) – Returns the path of every field a schema declares, given the schema’s id.
- Raises:
RobotoReadPlanExecutionException – With kind
plan-without-schema, when the plan projects every field of its schema, a partition has a scan task, and the plan names no schema.- Return type:
list[roboto.domain.topics.record.FieldPath]