---
sidebar:
  hidden: true
title: 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()`](/docs/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.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`](/docs/reference/python-sdk/roboto/experimental/topics/decode), combining a split partition in [`split_partition`](/docs/reference/python-sdk/roboto/experimental/topics/split_partition), and the output schema in that module and [`batch_transforms`](/docs/reference/python-sdk/roboto/experimental/topics/batch_transforms).

## Module Contents

### SchemaFieldPaths

```python
roboto.experimental.topics.plan_execution.SchemaFieldPaths
```

`from roboto.experimental.topics.plan_execution import SchemaFieldPaths`

[Source](https://github.com/roboto-ai/roboto-python-sdk/blob/main/src/roboto/experimental/topics/plan_execution.py#L48-L48)

Returns the path of every field a schema declares, given the schema's id.

### assign_scan_tasks()

```python
def roboto.experimental.topics.plan_execution.assign_scan_tasks(
    partition: roboto.experimental.topics.read_plan.ReadPlanPartition,
    leaf_most: collections.abc.Sequence[roboto.domain.topics.record.FieldPath],
) -> list[roboto.experimental.topics.decode.common.ScanTaskGroup]
```

`from roboto.experimental.topics.plan_execution import assign_scan_tasks`

[Source](https://github.com/roboto-ai/roboto-python-sdk/blob/main/src/roboto/experimental/topics/plan_execution.py#L116-L190)

Return the groups of `partition`'s scan tasks that supply the leaf-most paths `leaf_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.

**Parameters**

- **partition** (`roboto.experimental.topics.read_plan.ReadPlanPartition`)
- **leaf_most** (`collections.abc.Sequence[roboto.domain.topics.record.FieldPath]`)

**Raises**

- [`RobotoReadPlanExecutionException`](/reference/python-sdk/roboto/exceptions/domain#roboto.exceptions.domain.RobotoReadPlanExecutionException): With kind `inconsistent-scan-tasks-on-file`, when scan tasks on one file disagree on its format, transformations or topic name; with kind `projected-field-in-no-scan-task`, when no scan task contains a path of `leaf_most`.

**Returns**

- `list[roboto.experimental.topics.decode.common.ScanTaskGroup]`

### execute_read_plan()

```python
def roboto.experimental.topics.plan_execution.execute_read_plan(
    plan: roboto.experimental.topics.read_plan.ReadPlan,
    projected: collections.abc.Sequence[roboto.domain.topics.record.FieldPath],
    open_file_decoder: roboto.experimental.topics.decode.common.FileDecoderOpener,
) -> collections.abc.Generator[pyarrow.RecordBatch, None, None]
```

`from roboto.experimental.topics.plan_execution import execute_read_plan`

[Source](https://github.com/roboto-ai/roboto-python-sdk/blob/main/src/roboto/experimental/topics/plan_execution.py#L193-L271)

Decode the files a read plan names and yield the topic's rows as RecordBatches.

A read takes these steps:

1. Projected paths: `projected`, from [`projected_paths()`](/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.projected_paths).
2. Leaf-most paths ([`leaf_most_paths()`](/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.leaf_most_paths)).
3. Assign scan tasks ([`assign_scan_tasks()`](/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.assign_scan_tasks)).
4. Group by file ([`assign_scan_tasks()`](/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.assign_scan_tasks)).
5. Decode each file, through the decoders `open_file_decoder` opens.
6. Combine a split partition ([`SplitPartition`](/reference/python-sdk/roboto/experimental/topics/split_partition#roboto.experimental.topics.split_partition.SplitPartition)).
7. Output schema ([`topic_data_schema()`](/reference/python-sdk/roboto/experimental/topics/batch_transforms#roboto.experimental.topics.batch_transforms.topic_data_schema)).
8. Partitions: read in plan order, each one's schema checked against the first partition's when its files open, before any of its rows.

Each partition yields only the rows inside its own [`window`](/reference/python-sdk/roboto/experimental/topics/read_plan#roboto.experimental.topics.read_plan.ReadPlanPartition.window), which can be narrower than the plan's: a read scoped to a Session that holds a file over part of its time span returns that part of the file and nothing else.

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()`](/reference/python-sdk/roboto/experimental/topics/plan_execution#roboto.experimental.topics.plan_execution.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()`](/reference/python-sdk/roboto/experimental/topics/batch_transforms#roboto.experimental.topics.batch_transforms.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`](/reference/python-sdk/roboto/exceptions/domain#roboto.exceptions.domain.RobotoReadPlanExecutionException): With the kind of the step that refuses the plan: `projected-field-in-no-scan-task` or `inconsistent-scan-tasks-on-file` (steps 3 and 4); `field-split-inside-non-struct` or `scan-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 kinds `open_file_decoder` and its decoders raise, such as `unsupported-format`, `unsupported-timestamp`, `invalid-timestamp`, `field-not-in-file` and `data-range-not-in-file` (step 5).

**Returns**

- `collections.abc.Generator[pyarrow.RecordBatch, None, None]`

### leaf_most_paths()

```python
def roboto.experimental.topics.plan_execution.leaf_most_paths(
    paths: collections.abc.Iterable[roboto.domain.topics.record.FieldPath],
) -> list[roboto.domain.topics.record.FieldPath]
```

`from roboto.experimental.topics.plan_execution import leaf_most_paths`

[Source](https://github.com/roboto-ai/roboto-python-sdk/blob/main/src/roboto/experimental/topics/plan_execution.py#L100-L113)

Return the paths among `paths` that 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]`)

**Returns**

- `list[roboto.domain.topics.record.FieldPath]`

### projected_paths()

```python
def roboto.experimental.topics.plan_execution.projected_paths(
    plan: roboto.experimental.topics.read_plan.ReadPlan,
    schema_field_paths: SchemaFieldPaths,
) -> list[roboto.domain.topics.record.FieldPath]
```

`from roboto.experimental.topics.plan_execution import projected_paths`

[Source](https://github.com/roboto-ai/roboto-python-sdk/blob/main/src/roboto/experimental/topics/plan_execution.py#L72-L97)

Return the field paths `plan` projects.

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 through `schema_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`](/reference/python-sdk/roboto/exceptions/domain#roboto.exceptions.domain.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.

**Returns**

- `list[roboto.domain.topics.record.FieldPath]`
