ExArrow. Scanner
(ex_arrow v0.9.0)
View Source
Lazy scan of an ExArrow.Dataset with projection, partition pruning, and
filter pushdown.
Building a scanner (new/2 / ExArrow.Dataset.scanner/2) does no IO.
to_stream/1 starts an Agent-backed ExArrow.Stream (backend: :dataset)
that opens fragments on demand in path-sorted order.
Pushdown ladder
- Partition pruning — predicates on Hive keys are evaluated against
each fragment's
partition_values(no file open). - Parquet filters — remaining pushable predicates on data columns
become
Parquet.Reader:filters(row-group statistics). Hive keys are stripped from this AST because they are not columns in the file. - Residual — anything left runs through
ExArrow.Compute.filter/2after decode. Partition fields in a residual expression are bound to scalars for the current fragment.
Options
:columns— list of column names to project (Parquet pushdown / IPCCompute.project/2):filter—ExArrow.Compute.Expression.t(), legacy Parquet filter tuple, ornil:batch_size— accepted for API stability; reserved in 0.9 (batches follow Parquet row-group sizing)
Typical usage
alias ExArrow.Compute.Expression, as: E
{:ok, dataset} =
ExArrow.Dataset.open("/data/events",
partitioning: {:hive, schema: [{"year", :int32}, {"month", :int32}]}
)
{:ok, scanner} =
ExArrow.Dataset.scanner(dataset,
columns: ["id", "amount"],
filter:
E.and_(
E.gte(E.field("year"), E.scalar(2026)),
E.gt(E.field("amount"), E.scalar(0.0))
)
)
# Optional: preview how many fragments survive partition prune (no IO)
preview = ExArrow.Scanner.stats(scanner)
{:ok, stream} = ExArrow.Scanner.to_stream(scanner)
batches = Enum.to_list(stream)
stats = ExArrow.Scanner.stats(stream)
:ok = ExArrow.Stream.close(stream)Early Enum.take/2 does not open later fragments. Call
ExArrow.Stream.close/1 when abandoning a partially consumed scan from a
long-lived process.
Summary
Functions
Build a lazy scanner over dataset.
Return scan statistics.
Start scanning: partition-prune, then return an ExArrow.Stream.
Types
@type stats() :: %{ fragments_discovered: non_neg_integer(), fragments_pruned_partition: non_neg_integer(), fragments_selected: non_neg_integer(), fragments_scanned: non_neg_integer(), row_groups_selected: non_neg_integer(), row_groups_skipped: non_neg_integer(), rows_emitted: non_neg_integer() }
Cumulative scan statistics.
Keys
:fragments_discovered— fragments on the Dataset before prune:fragments_pruned_partition— dropped by Hive-key evaluation:fragments_selected— kept after partition prune:fragments_scanned— actually opened duringto_stream/1:row_groups_selected/:row_groups_skipped— sums ofParquet.Reader.read_stats/1across opened Parquet fragments:rows_emitted— rows yielded after residual filter (empty batches are skipped and do not count)
@type t() :: %ExArrow.Scanner{ batch_size: pos_integer() | nil, columns: [String.t()] | nil, dataset: ExArrow.Dataset.t(), filter: ExArrow.Compute.Expression.t() | tuple() | nil, partition_keys: [String.t()] }
A lazy scan plan over a Dataset.
Fields
:dataset— sourceExArrow.Dataset.t():columns— projection list, ornilfor all columns:filter— Expression, legacy filter tuple, ornil:batch_size— reserved positive integer, ornil:partition_keys— Hive column names used when compiling filters (derived from the dataset partitioning schema)
Functions
@spec new( ExArrow.Dataset.t(), keyword() ) :: {:ok, t()} | {:error, String.t()}
Build a lazy scanner over dataset.
Performs no IO. Prefer ExArrow.Dataset.scanner/2, which delegates here.
Parameters
dataset— anExArrow.Dataset.t()opts— keyword list::columns— non-empty[String.t()]to project, or omit/nil:filter—Expression.t(), legacy tuple, or omit/nil:batch_size— positive integer (reserved; validated only)
Returns
{:ok, scanner}when options validate against the dataset schema (Expression fields may include Hive partition columns){:error, message}for bad options or filter validation failures
Examples
{:ok, scanner} = ExArrow.Scanner.new(dataset, columns: ["id"])
{:ok, scanner} = ExArrow.Scanner.new(dataset, filter: {:gt, "amount", 0.0})
@spec stats(t() | ExArrow.Stream.t()) :: stats() | {:error, String.t()}
Return scan statistics.
Parameters
scanner_or_stream— either:- an
ExArrow.Scanner.t()— preview after partition prune only (fragments_scanned, row-group, androws_emittedare0) - an
ExArrow.Stream.t()withbackend: :dataset— live or post-scan aggregates from the Agent
- an
Returns
A stats/0 map for an open scanner or live stream. For a :dataset
stream whose Agent has already been stopped via ExArrow.Stream.close/1,
returns {:error, "stream is closed"} instead of raising.
Examples
preview = ExArrow.Scanner.stats(scanner)
preview.fragments_pruned_partition
{:ok, stream} = ExArrow.Scanner.to_stream(scanner)
_ = Enum.to_list(stream)
stats = ExArrow.Scanner.stats(stream)
stats.row_groups_skipped
stats.rows_emitted
:ok = ExArrow.Stream.close(stream)
{:error, "stream is closed"} = ExArrow.Scanner.stats(stream)
@spec to_stream(t()) :: {:ok, ExArrow.Stream.t()} | {:error, String.t()}
Start scanning: partition-prune, then return an ExArrow.Stream.
The stream has backend: :dataset and implements Enumerable. Fragments
open lazily; EOF is idempotent (next/1 returns nil forever after
exhaustion).
Parameters
scanner— fromnew/2/Dataset.scanner/2
Returns
{:ok, stream}— Agent-backed stream; callExArrow.Stream.close/1when done or when abandoning early{:error, message}— filter compile / Parquet opts validation failure
Examples
{:ok, stream} = ExArrow.Scanner.to_stream(scanner)
batch = ExArrow.Stream.next(stream)
batches = Enum.to_list(stream)
:ok = ExArrow.Stream.close(stream)Telemetry: emits [:ex_arrow, :dataset, :scan, :start] here and
[:ex_arrow, :dataset, :scan, :stop] when the scan finishes or is closed.
Per-batch events use [:ex_arrow, :stream, :batch] with
source: {:dataset, current_fragment_path}.