Python API Reference¶
This page is a reference for the rypipe Python API. For a tutorial, see the Tutorial.
rypipe.read()¶
Read a file into a pyarrow.Table using a registered adapter.
rypipe.read(
path, # str | PathLike: file path
*,
format=None, # str | None: adapter name
adapter=None, # object with read() method
**kwargs, # forwarded to the adapter
) -> pyarrow.Table
Parameters:
| Parameter | Type | Default | Description |
|---|---|---|---|
path |
str \| PathLike |
(required) | Path to the input file. |
format |
str \| None |
None |
Adapter name. Inferred from extension when omitted. |
adapter |
Any \| None |
None |
Adapter object. Overrides format. |
**kwargs |
: | : | Forwarded to the adapter's read() method. |
Returns: pyarrow.Table
Raises: rypipe.RypipeError if no adapter is registered for the format.
import rypipe
import crxml # registers the crxml adapter
# Format inferred from extension
table = rypipe.read("report.xml", row_tag="Details")
# Format specified explicitly
table = rypipe.read("data.txt", format="crxml")
# Adapter passed directly
from crxml import CrystalXMLAdapter
table = rypipe.read("report.xml", adapter=CrystalXMLAdapter(), row_tag="Details")
rypipe.read_par()¶
Read a file in parallel using a registered adapter.
rypipe.read_stream()¶
Read a file with bounded memory using a registered adapter.
rypipe.read_stream(
path,
*,
memory="64MiB", # int | str: memory budget
**kwargs,
) -> pyarrow.Table
rypipe.read_batches()¶
Read a file and yield pyarrow.RecordBatch objects incrementally.
rypipe.read_batches(
path,
*,
memory="64MiB", # int | str: memory budget
batch_size=None, # int | None: rows per batch
**kwargs,
) -> Iterator[pyarrow.RecordBatch]
rypipe.iter_record_batches()¶
Stream a file into Arrow RecordBatch objects with constant memory.
rypipe.iter_record_batches(
path,
*,
format=None, # str | None: adapter name
adapter=None, # object with read() method
memory="64MiB", # int | str: memory budget
batch_size=None, # int | None: rows per batch
**kwargs,
) -> Iterator[pyarrow.RecordBatch]
rypipe.register_adapter()¶
Register a format adapter with rypipe.
rypipe.register_adapter(
name, # str: adapter name
adapter, # object with read() method
extensions=None, # Iterable[str] | None: file extensions
) -> None
After registration, rypipe.read("file.ext") auto-detects the extension.
Source¶
Abstract base class for row-oriented file sources.
class Source(ABC):
def __init__(
self,
path, # str | Path
*,
field_mapping=None, # dict[str, str]
drop_fields=None, # list[str]
filter=None, # dict | None
field_types=None, # dict[str, str]
dictionary_columns=None, # list[str]
schema=None, # list[str]
auto_dict=False, # bool
use_mmap=True, # bool
batch_size=1024, # int
)
Abstract method:
Public methods:
| Method | Returns | Description |
|---|---|---|
.to_arrow() |
pyarrow.Table |
Parse and cache the table. |
.to_pandas(dtype_backend="pyarrow") |
pd.DataFrame |
Convert to pandas. |
.to_dataframe(dtype_backend="pyarrow") |
pd.DataFrame |
Alias for .to_pandas(). |
.to_polars() |
pl.DataFrame |
Convert to Polars. |
.to_parquet(path, **kwargs) |
None |
Write to Parquet. |
.schema() |
list[str] |
Column names from first row. |
.clear_cache() |
None |
Drop cached table. |
.iter_arrow_batches(batch_size=None) |
Iterator[RecordBatch] |
Yield batches. |
.iter_record_batches(memory="64MiB", batch_size=None) |
Iterator[RecordBatch] |
Stream batches. |
.__iter__() |
Iterator[dict] |
Iterate rows as dicts. |
.__or__(stage) |
Pipeline |
Pipe operator for stages. |
Adapter¶
Convenience base class. Subclasses implement read() instead of
_read_arrow().
class Adapter(Source):
def read(self, path: str, **kwargs) -> pyarrow.Table:
raise NotImplementedError
Pipeline¶
A chain of stages applied to a Source.
class Pipeline:
def __or__(self, stage) -> Pipeline:
"""Append a stage and return a new Pipeline."""
def __iter__(self) -> Iterator[dict]:
"""Iterate rows as dicts."""
def iter_arrow_batches(self, batch_size=None) -> Iterator[RecordBatch]:
"""Yield Arrow RecordBatch objects."""
def iter_record_batches(self, memory="64MiB", batch_size=None) -> Iterator[RecordBatch]:
"""Stream batches with constant memory."""
Stages¶
Pipeline stages transform streams of dicts. Import from the adapter package.
RenameFields¶
Rename columns. Fields not in the mapping pass through unchanged.
DropFields¶
Remove columns. Passing a bare string raises TypeError.
CastTypes¶
Cast column values. Supported callables: int, float, str, bool.
FilterRows¶
FilterRows(
predicate=None, # Callable: arbitrary filter
*,
field=None, # str: column name (constant filter)
op=None, # str: operator
value=None, # str: value (constant filter)
field_a=None, # str: left column (comparison)
field_b=None, # str: right column (comparison)
)
Constant filter operators: ==, eq, !=, ne
Comparison operators: >, <, >=, <=, ==, !=, gt, lt,
ge, le, eq, ne
FilterRowsAny¶
Keep rows matching any filter (OR).
FilterRowsAll¶
Keep rows matching all filters (AND).
FilterRowsNot¶
Negate a filter.
Sinks¶
Standalone functions for materializing pipeline results.
| Function | Returns | Description |
|---|---|---|
rypipe.collect(pipeline) |
list[dict] |
Collect all rows. |
rypipe.to_arrow(pipeline) |
pyarrow.Table |
Materialize to table. |
rypipe.to_pandas(pipeline) |
pd.DataFrame |
Materialize to pandas. |
rypipe.to_polars(pipeline) |
pl.DataFrame |
Materialize to Polars. |
rypipe.to_csv(pipeline, path, ...) |
None |
Write to CSV. |
rypipe.to_parquet(pipeline, path, ...) |
None |
Write to Parquet. |
to_csv()¶
rypipe.to_csv(
pipeline, # Iterable[dict]
path, # str | Path
encoding="utf-8", # str
delimiter=",", # str
fieldnames=None, # list[str] | None
) -> None
Exceptions¶
| Exception | Parent | Meaning |
|---|---|---|
rypipe.RypipeError |
RuntimeError |
General API error. |
rypipe.ParseError |
Exception |
File could not be parsed. |
rypipe.XmlError |
ParseError |
XML-specific parse error. |
rypipe.PlanError |
Exception |
Invalid plan kwargs. |
rypipe.MergeError |
Exception |
Schema mismatch between chunks. |