Pushdown fusion¶
rypipe splits ingestion into two layers:
- Adapter layer: reads bytes, splits them into records, and emits
Valuerows. - Engine layer: builds typed Arrow arrays, applies filters, and exports record batches.
Pushdown fusion is the process by which the Python Pipeline rewrites a chain of lightweight stages into a single ExecutionPlan. The engine then applies rename, drop, filter, and cast while it parses each row, instead of materializing a full table and running Python transforms afterward.
A fused pipeline¶
# Stages are imported from the adapter package; they are the same
# classes the framework's stage protocol defines.
from crxml import CrystalXMLSource, RenameFields, DropFields, FilterRows, CastTypes, to_pandas, col
source = CrystalXMLSource("report.xml", row_tag="Details")
df = to_pandas(
source
| RenameFields({"old_name": "new_name"})
| DropFields(["internal_id"])
| FilterRows(field="status", op="==", value="active")
| CastTypes({"amount": float})
)
# Expression predicates are fusable too:
df = to_pandas(
source
| FilterRows((col("amount") > 100) & (col("status") == "active"))
)
Expression predicates (col("x") > 100, col("s").startswith("A"), col("s").matches(r"^ERR"), col("n").between(1, 10), combined with &, |, ~) build the same filter spec dicts the engine fuses, with no bytecode analysis. Plain lambdas and other callables still work but always run in Python; see expression filters.
When the pipeline is materialized, the stage list is collapsed into one plan. The Rust parser:
- renames
old_nametonew_nameas fields arrive; - skips
internal_identirely (it is not allocated); - drops rows whose
statusis notactivebefore they leave the builder; - casts
amounttofloat64once, during parse.
Rows that fail the filter are never materialized; dropped columns are never allocated; and casts happen once inside Rust instead of twice in Python and Rust.
Fusion pipeline flow¶
Python pipeline Rust engine
───────────── ───────────
source | RenameFields | DropFields ExecutionPlan
| FilterRows | CastTypes ├── field_map: {old_name → new_name}
| .to_arrow() ├── drop_fields: {internal_id}
├── field_types: {amount → Float64}
│ ├── filter: Equal(status, "active")
▼ └── schema_order: []
plan_split() │
│ ▼
plan_overrides = { TableBuilder::push_field()
field_mapping, │
drop_fields, ┌──────┴──────┐
field_types, │resolve name │
filter │via field_map│
} └──────┬──────┘
│ │
▼ ┌────────┴────────┐
source._read_arrow(plan) │ drop? → skip │
│ │ cast → typed │
▼ │ filter → reject │
Rust parse + export │ push → column │
│ └───────┬─────────┘
▼ │
pyarrow.Table finish_row()
│ └── null-fill, filter check
▼
to_arrow() result
Steps 1-4 happen once per row inside the Rust parse loop. No Python objects are created for fused stages.
The ExecutionPlan fields¶
pub struct ExecutionPlan {
pub field_map: HashMap<String, String>, // rename
pub drop_fields: HashSet<String>, // drop
pub field_types: HashMap<String, FieldType>, // cast
pub dictionary_columns: HashSet<String>, // explicit dict encoding
pub filter: Option<FilterPredicate>, // per-row (all kinds)
pub schema_order: Vec<String>, // output column order
pub auto_dict: bool, // auto-dict upgrade
pub dict_threshold: Option<f64>, // auto-dict ratio (default 0.05)
pub dict_max_size: Option<usize>, // auto-dict max entries (default 256)
pub strict_types: bool, // abort on non-null parse failures
pub max_split_chunks: Option<usize>, // bounded-streaming chunk cap
pub min_chunk_bytes: Option<usize>, // split chunk-size floor
pub observer: Option<Arc<dyn RowObserver>>, // per-row observer hooks
}
The plan is built by _build_plan_kwargs() on the Python side and consumed by TableBuilder on the Rust side. The Python-facing kwargs accepted by rypipe-python entry points are: field_mapping, drop_fields, filter, field_types, dictionary_columns, schema (mapped to schema_order), auto_dict, auto_dict_threshold (mapped to dict_threshold), auto_dict_max_size (mapped to dict_max_size), and min_chunk_bytes.
Field resolution order¶
Inside TableBuilder, every emitted field goes through this pipeline:
field_maprenames the raw field name.drop_fieldschecks the resolved name; if dropped, the field is ignored.field_types/dictionary_columnschooses the storage type.strict_typesvalidates the value parses as the declared type (records violation forfinishto report).observerfireson_put_field(ifplan.observeris set).- Value is pushed into the column; dirty bit is set.
The per-row filter is evaluated later, in finish_row, after all fields
are pushed and missing columns are null-filled.
This order matters. A filter runs on the resolved name, so it must be written in post-rename terms. A cast type is attached to the resolved name as well.
What is fusable¶
Fusable stages implement _plan_kwargs() and merge cleanly into an ExecutionPlan:
| Stage | Plan field | Notes |
|---|---|---|
RenameFields |
field_map |
A contiguous prefix merges into one map; repeated or out-of-order renames use the fallback path. |
DropFields |
drop_fields |
Merges as a set union. |
CastTypes |
field_types |
A contiguous prefix fuses; conflicting or mixed cast stages use the fallback path. |
FilterRows predicate |
filter |
Keyword form (field/op/value or field_a/op/field_b, including op="regex"), is_null, is_type, or an expression predicate built with the adapter's re-exported col (comparisons, startswith, endswith, contains, matches, between, isin, not_in, compound &/|/~); all are evaluated per-row during parse. |
FilterRowsAny / FilterRowsAll / FilterRowsNot |
filter |
And, Or, Not trees built from the same leaf shapes; evaluated per-row with short circuiting; fully fusable. |
ObservedStage (or any stage returning observer hooks) |
observer |
Hook dicts merge per-hook; callables for the same hook chain in stage order. |
CastTypes treats int, float, bool, date, datetime, and Decimal as
column types. Cached and direct execution reuse Rust's string conversions,
so "false" becomes False, null values stay null, and decimal truncation
matches a fresh read. Already typed Arrow columns use Arrow casts. Other
callables, including str, run as Python functions. Mixed stages run outside
the Rust plan without dropping their custom conversions.
FilterRows is fusable when it uses a keyword-form predicate (field,
op, value or field_a, op, field_b), or an expression predicate
built with the adapter's re-exported col (see expression filters).
FilterRowsAny, FilterRowsAll, and FilterRowsNot are also fusable;
they build And, Or, Not trees from the same leaves. All are evaluated per-row during parsing with native-typed comparison and numeric promotion; mismatched types or nulls fail the row, with Not flipping the result. Chaining FilterRows stages is an implicit And (see plan_split).
What is not fusable¶
Non-fusable stages still work, but they run over the Arrow table after the engine finishes:
- Raw Python callables used as standalone pipeline stages.
- Stateful transforms such as window or aggregate stages.
FilterRowswrapping a plain lambda or named function (Python fallback only).- Custom stages that do not implement
_plan_kwargs().
When non-fusable stages are present, only the valid contiguous fusable prefix is
sent to the Rust plan. Out-of-order stages, repeated renames, and mixed
CastTypes chains stop fusion at that point; remaining stages run over the
materialized Arrow batches in Python.
Inspecting plan_overrides in an adapter¶
Adapters that subclass rypipe.Source receive fused plan kwargs through _read_arrow:
class MySource(Source):
def _read_arrow(self, plan_overrides=None):
plan = self._build_plan_kwargs()
if plan_overrides:
plan.update(plan_overrides)
print(plan)
# {
# "use_mmap": True,
# "auto_dict": False,
# "field_mapping": {"old_name": "new_name"},
# "drop_fields": ["internal_id"],
# "filter": {"field": "status", "op": "==", "value": "active"},
# "field_types": {"amount": "float64"},
# }
return my_rust_read(str(self._path), **plan)
The Rust side merges these kwargs into an ExecutionPlan with execution_plan_from_kwargs from rypipe-python, or constructs the plan manually:
use rypipe_core::{ExecutionPlan, FieldType};
let plan = ExecutionPlan::new()
.rename("old_name", "new_name")
.drop("internal_id")
.type_as("amount", FieldType::Float64)
.filter_eq("status", "active");
If an adapter ignores plan_overrides, fused stage transformations are absent
from the Rust plan. A reader that rejects the unexpected keywords raises
TypeError; forwarding the overrides preserves the requested transformations.
Order of operations across the pipeline¶
Fusable stages commute in the plan, but the engine applies them in a fixed order:
raw field name
|
v
rename (field_map)
|
v
drop check (drop_fields)
|
v
type selection (field_types / dictionary_columns)
|
v
strict_types validation (records parse violations)
|
v
observer notification (on_put_field)
|
v
push value + set dirty bit
|
v
[end_row: null-fill missing columns, then per-row filter]
|
v
finish: sort by schema_order, auto-dict, Arrow export
Because drop happens before type selection, you cannot cast a dropped field.
Because the filter runs after all fields are pushed (in finish_row),
rejected rows consume no Arrow storage but their column values are
temporarily allocated then popped.
When fusion does not help¶
Fusion is not free if the adapter cannot act on the plan. An adapter that always parses every field into Python objects and then builds Arrow will not benefit; the engine must receive fields through put_field and honor wants(). If the adapter is a thin wrapper around a library that returns full Python dicts, fusion only removes a small amount of Python overhead.
Fusion also cannot help when the workload is dominated by I/O. If the file is on a slow network share, reducing CPU work may not change wall-clock time. Profile first.
Summary¶
- Fuse
RenameFields,DropFields,CastTypes, andFilterRowspredicates (constant and column-to-column) by implementing_read_arrow(plan_overrides=...). - Inspect
plan_overridesto confirm that stages reach the Rust parser. - Non-fusable stages run after the engine; keep them out of the hot path when throughput matters.