Sinks¶
Note
The sinks shown here are available in most adapters but may have different options. Check your adapter's docs.
Sinks materialize pipeline results into tables, DataFrames, or files. You can use them as methods on a Source or as standalone functions on a Pipeline.
Source methods¶
Every Source has built-in sink methods:
to_arrow()¶
Returns a pyarrow.Table. This is the default materialization:
to_pandas()¶
Returns a pandas DataFrame with PyArrow-backed dtypes by default:
You can disable PyArrow backing with dtype_backend="numpy":
to_dataframe()¶
Alias for to_pandas():
to_polars()¶
Returns a Polars DataFrame:
to_parquet()¶
Writes the table to a Parquet file:
src.to_parquet("output.parquet")
# Pass additional pyarrow.parquet options
src.to_parquet("output.parquet", compression="snappy")
clear_cache()¶
Drops the cached Arrow table to free memory:
Pipeline functions¶
When working with a Pipeline (the result of src | stage), use the
standalone sink functions from the adapter:
from crxml import CrystalXMLSource, FilterRows, collect
src = CrystalXMLSource("report.xml", row_tag="Details")
pipeline = src | FilterRows(field="status", op="==", value="active")
collect()¶
Collects all rows into a list of dicts:
rows = rypipe.collect(pipeline)
print(rows[0])
# {"name": "Alice", "amount": 150.0, "status": "active"}
to_arrow() (function)¶
Materializes a pipeline to a pyarrow.Table:
to_pandas() (function)¶
Materializes a pipeline to a pandas DataFrame:
to_polars() (function)¶
Materializes a pipeline to a Polars DataFrame:
to_csv() (function)¶
Writes pipeline results to a CSV file:
rypipe.to_csv(pipeline, "output.csv")
# Custom delimiter and encoding
rypipe.to_csv(pipeline, "output.tsv", delimiter="\t", encoding="utf-8")
Parameters:
pipeline: iterable of dicts.path: output file path.encoding: file encoding (default:"utf-8").delimiter: column delimiter (default:",").fieldnames: optional list of column names. If omitted, uses the keys from the first row.
to_parquet() (function)¶
Writes pipeline results to a Parquet file:
Which sink should I use?¶
| Goal | Method |
|---|---|
| Get a PyArrow table | .to_arrow() or rypipe.to_arrow() |
| Get a pandas DataFrame | .to_pandas() or rypipe.to_pandas() |
| Get a Polars DataFrame | .to_polars() or rypipe.to_polars() |
| Write to Parquet | .to_parquet(path) or rypipe.to_parquet(pipeline, path) |
| Write to CSV | rypipe.to_csv(pipeline, path) |
| Get a list of dicts | rypipe.collect(pipeline) |
Tip
When you have a Source, prefer the Source methods (.to_pandas(), etc.)
over the standalone functions. Source methods reuse the cached table and
avoid re-parsing.
Recap¶
- Source methods:
.to_arrow(),.to_pandas(),.to_polars(),.to_parquet(),.clear_cache(). - Standalone functions:
rypipe.collect(),rypipe.to_arrow(),rypipe.to_pandas(),rypipe.to_polars(),rypipe.to_csv(),rypipe.to_parquet(). - Source methods reuse the cached table. Standalone functions re-parse if the pipeline hasn't been materialized yet.
Next: Streaming: processing large files with bounded memory.