Skip to content

Pushdown fusion

rypipe splits ingestion into two layers:

  1. Adapter layer: reads bytes, splits them into records, and emits Value rows.
  2. 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_name to new_name as fields arrive;
  • skips internal_id entirely (it is not allocated);
  • drops rows whose status is not active before they leave the builder;
  • casts amount to float64 once, 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:

  1. field_map renames the raw field name.
  2. drop_fields checks the resolved name; if dropped, the field is ignored.
  3. field_types / dictionary_columns chooses the storage type.
  4. strict_types validates the value parses as the declared type (records violation for finish to report).
  5. observer fires on_put_field (if plan.observer is set).
  6. 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.
  • FilterRows wrapping 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, and FilterRows predicates (constant and column-to-column) by implementing _read_arrow(plan_overrides=...).
  • Inspect plan_overrides to confirm that stages reach the Rust parser.
  • Non-fusable stages run after the engine; keep them out of the hot path when throughput matters.