Streaming¶
Note
Streaming options (memory, batch_size) are available in most adapters
but may have different defaults. Check your adapter's docs.
When processing files larger than available memory, use streaming to process data in bounded chunks.
The problem¶
By default, .to_arrow() parses the entire file into memory at once.
This works for files up to several GB on a machine with enough RAM,
but fails for larger files:
from crxml import CrystalXMLSource
# This loads the entire file into memory
source = CrystalXMLSource("huge_report.xml", row_tag="Details")
table = source.to_arrow() # may OOM
Streaming with iter_record_batches¶
The iter_record_batches() method yields pyarrow.RecordBatch objects one
at a time, processing only a bounded amount of data at once:
from crxml import CrystalXMLSource
src = CrystalXMLSource("huge_report.xml", row_tag="Details")
# Process in chunks of ~64 MiB
for batch in src.iter_record_batches(memory="64MiB"):
# Each batch is a pyarrow.RecordBatch
print(f"Processing {batch.num_rows} rows")
process(batch)
How it works¶
- rypipe reads a chunk of the file into memory (bounded by
memory). - Parses the chunk into a
RecordBatch. - Yields the batch to your code.
- Drops the batch from memory after your code returns.
- Repeats for the next chunk.
Peak memory is memory + one batch + export buffer. A 10 GB file with
memory="256MiB" uses at most ~300 MB of parsing memory at any time.
Memory parameter¶
The memory parameter accepts a string or integer:
# String formats
for batch in src.iter_record_batches(memory="64MiB"):
...
for batch in src.iter_record_batches(memory="256MB"):
...
# Integer (bytes)
for batch in src.iter_record_batches(memory=67_108_864): # 64 MiB
...
Supported units: B, KB, MB, GB, TB, KiB, MiB, GiB, TiB.
Batch size¶
Control the number of rows per batch with batch_size:
# Smaller batches = lower memory, more Python overhead
for batch in src.iter_record_batches(memory="64MiB", batch_size=1000):
process(batch)
# Larger batches = higher memory, less overhead
for batch in src.iter_record_batches(memory="256MiB", batch_size=100_000):
process(batch)
When batch_size is None (the default), rypipe sizes batches based
on the memory budget and estimated row size.
Streaming with pipelines¶
Pipelines also support streaming:
from crxml import CrystalXMLSource
from crxml import CastTypes, FilterRows
src = CrystalXMLSource("huge_report.xml", row_tag="Details")
pipeline = (
src
| CastTypes({"amount": float})
| FilterRows(field="status", op="==", value="active")
)
for batch in pipeline.iter_record_batches(memory="64MiB"):
process(batch)
Note
When all stages are fusable, rypipe pushes the entire pipeline into the streaming parse loop. Non-fusable stages run after each batch is parsed.
Writing to Parquet¶
A common pattern is streaming a large file into a Parquet writer:
import pyarrow.parquet as pq
from crxml import CrystalXMLSource
from crxml import CastTypes
src = CrystalXMLSource("huge_report.xml", row_tag="Details")
pipeline = src | CastTypes({"amount": float})
# Write in streaming mode
writer = pq.ParquetWriter("output.parquet", schema=None)
for batch in pipeline.iter_record_batches(memory="64MiB"):
if writer.schema is None:
# First batch: initialize the writer with the schema
writer = pq.ParquetWriter(
"output.parquet",
batch.schema,
)
writer.write_batch(batch)
writer.close()
Streaming from a Source¶
You can also stream directly from any Source:
from crxml import CrystalXMLSource
src = CrystalXMLSource("huge_report.xml", row_tag="Details")
for batch in src.iter_record_batches(memory="64MiB"):
process(batch)
Performance¶
Streaming has slightly lower throughput than full-table parsing because batches are processed one at a time rather than in parallel. Typical numbers for the crxml adapter:
| Mode | Throughput | Peak memory |
|---|---|---|
| Parallel (default) | ~4 GB/s | File size |
| Single-thread | ~1 GB/s | File size |
| Streaming (64 MiB) | ~500 MB/s | ~64 MiB |
Tip
Use streaming when your file is larger than ~50% of available RAM. For smaller files, the default parallel mode is faster.
Recap¶
- Use
.iter_record_batches(memory="64MiB")for bounded-memory processing. - Peak memory is
memory+ one batch + export buffer. - Control batch size with
batch_sizefor tuning memory vs overhead. - Streaming works with pipelines: fusable stages run in the parse loop.
- Use
rypipe.iter_record_batches()for module-level streaming.
Next: Configuration: all available options.