Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,47 @@ All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [0.9.0] - 2026-09-13

### Added

- **`ExArrow.Dataset` / `ExArrow.Dataset.Fragment`**: discover Parquet or IPC
trees from a directory, single file, glob, or explicit path list via
`Dataset.open/2`. Hive partitioning parses typed keys from path segments;
each fragment carries `path`, `format`, `size`, and `partition_values`.
Schema is resolved from the first fragment footer without decoding pages.
- **`ExArrow.FileSystem`**: behaviour with `Local` (OS) and `Memory`
(in-process tree) backends for list/glob/exists used by Dataset discovery.
- **`ExArrow.Scanner`**: lazy scan plan from `Dataset.scanner/2` with
`:columns`, `:filter`, and `:batch_size`. `Scanner.to_stream/1` yields an
Agent-backed `:dataset` stream; `Scanner.stats/1` reports exact fragment
and row-group prune counts (preview before IO, or live after scan). Closed
streams return `{:error, "stream is closed"}` instead of exiting.
- **`ExArrow.Compute.Expression`**: analyzable filter AST with `field/1`,
`scalar/1`, comparisons, `and_/2`, `or_/2`, `not_/1`, plus `validate/2`,
`to_string/1`, and `to_parquet_filters/1` (pushable subset vs residual).
Temporal scalars (`Date`, `NaiveDateTime`, `DateTime`) are supported for
residual evaluation.
- **Residual `Compute.filter/2`**: new NIF evaluates an Expression to a
boolean mask and filters a `RecordBatch` (also via `Batch.filter/2`).
Integer→float casts reject values that are not exactly representable.
- **`RecordBatch.from_lists/1` and `from_map/1`**: ergonomic constructors
for non-nullable columnar batches from Elixir lists/maps.
- **Docs / fixtures**: Datasets guide (`guides/11_datasets.md`), Livebook
`06_datasets.livemd`, PyArrow Hive fixture with exact Scanner prune
stats, and `bench/dataset_scan_bench.exs` (labeled pushdown ladder).

### Changed

- **Native stack**: arrow-rs / parquet / arrow-flight **59.3.0** (from 56);
tonic **0.14**; adbc_core / adbc_driver_manager **0.24**.
- **Parquet `:filters`**: continues to accept the v0.8.0 tuple AST; an
`Expression` is accepted only when fully Parquet-pushable (use Scanner
when a residual remains).
- Cross-fragment Scanner schema checks compare **name and type** (not names
alone). Dataset fragment list building uses linear prepend/reverse
accumulation.

## [0.8.0] - 2026-08-21

### Added
Expand Down
99 changes: 95 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,12 @@

Native Apache Arrow for the BEAM: IPC streaming, Arrow Flight, Arrow Flight SQL, ADBC database bindings, and Arrow-native pipelines. Column data lives in Rust buffers; Elixir holds lightweight opaque handles. Precompiled NIFs for Linux, macOS, and Windows — no Rust required to use.

> **v0.9.0 — Dataset and Scanner.** Discover Hive-partitioned Parquet/IPC trees,
> filter with `ExArrow.Compute.Expression`, and scan through partition prune →
> row-group pushdown → residual `Compute.filter/2`. See
> [Datasets and scanners](#datasets-and-scanners) and the
> [Datasets guide](https://ex-arrow.hexdocs.pm/guides/11_datasets.html).
>
> **v0.8.0 — Larger-than-memory Parquet.** Column/predicate/row-group pushdown,
> write compression options, footer metadata, and multi-file directory streams.
> See [Parquet: read and write](#parquet-read-and-write) and the
Expand All @@ -31,6 +37,7 @@ Native Apache Arrow for the BEAM: IPC streaming, Arrow Flight, Arrow Flight SQL,
- [Requirements](#requirements)
- [Installation](#installation)
- [Quick start](#quick-start)
- [What's changed in v0.9.0](#whats-changed-in-v090)
- [What's changed in v0.8.0](#whats-changed-in-v080)
- [What's changed in v0.7.0](#whats-changed-in-v070)
- [Livebook tutorials](#livebook-tutorials)
Expand All @@ -55,6 +62,7 @@ Native Apache Arrow for the BEAM: IPC streaming, Arrow Flight, Arrow Flight SQL,
- [Shipped (v0.5.0)](#shipped-v050)
- [Shipped (v0.7.0)](#shipped-v070)
- [Shipped (v0.8.0)](#shipped-v080)
- [Shipped (v0.9.0)](#shipped-v090)
- [FAQ](#faq)
- [License](#license)

Expand Down Expand Up @@ -215,7 +223,7 @@ Add the dependency:

```elixir
def deps do
[{:ex_arrow, "~> 0.8"}]
[{:ex_arrow, "~> 0.9"}]
end
```

Expand Down Expand Up @@ -243,11 +251,11 @@ For **path dependencies** in Livebook (`Mix.install`), open notebooks from
is detected) or use the Hex package:

```elixir
Mix.install([{:ex_arrow, "~> 0.8.0"}, {:rustler, "~> 0.36", optional: true}])
Mix.install([{:ex_arrow, "~> 0.9.0"}, {:rustler, "~> 0.36", optional: true}])
```

Alternatively, use the published Hex package so the precompiled NIF is used
and no Rust is needed: `Mix.install([{:ex_arrow, "~> 0.8.0"}])`.
and no Rust is needed: `Mix.install([{:ex_arrow, "~> 0.9.0"}])`.

---

Expand Down Expand Up @@ -329,6 +337,39 @@ batch = ExArrow.Stream.next(stream)

---

## What's changed in v0.9.0

v0.9.0 adds a **Dataset / Scanner** layer on top of v0.8.0 Parquet pushdown:
discover partitioned trees, compile analyzable Expression filters, and scan
with a three-level pushdown ladder (Hive partition prune → Parquet row-group
pushdown → residual `Compute.filter/2`).

### Dataset, Scanner, Expression

| API | Purpose |
|-----|---------|
| `ExArrow.Dataset.open/2` | Open a directory, file, glob, or path list; Hive `partition_values` |
| `ExArrow.Dataset.Fragment` | One file fragment with path, format, size, partition map |
| `ExArrow.Dataset.scanner/2` | Lazy scan plan: `:columns`, `:filter`, `:batch_size` |
| `ExArrow.Scanner.to_stream/1` | Agent-backed `:dataset` stream; `Stream.close/1` for early abandon |
| `ExArrow.Scanner.stats/1` | Exact fragment / row-group prune counts (preview or live) |
| `ExArrow.Compute.Expression` | Builders, `validate/2`, `to_parquet_filters/1` (pushable vs residual) |
| `ExArrow.Compute.filter/2` | Residual evaluation of an Expression against a RecordBatch |
| `ExArrow.FileSystem` | `Local` and `Memory` discovery backends |
| `RecordBatch.from_lists/1`, `from_map/1` | Ergonomic batch construction |

Native stack: arrow-rs / parquet / arrow-flight **59.3.0** (from 56).

### Docs, Livebook, bench

- Datasets Livebook (`livebook/06_datasets.livemd`) and
[Datasets guide](guides/11_datasets.md)
- Pushdown-ladder timing helper: `bench/dataset_scan_bench.exs`

See [CHANGELOG.md](CHANGELOG.md) for the full list.

---

## What's changed in v0.8.0

v0.8.0 focuses on **larger-than-memory Parquet**: read only the columns, row
Expand Down Expand Up @@ -415,6 +456,7 @@ without `:telemetry`, the Flow/GenStage/Broadway modules return
- [08 Arrow and GenStage](guides/08_arrow_and_genstage.md)
- [09 Arrow and Broadway](guides/09_arrow_and_broadway.md)
- [10 Arrow pipeline patterns](guides/10_arrow_pipeline_patterns.md)
- [11 Datasets and scanners](guides/11_datasets.md)

### New benchmarks

Expand All @@ -435,8 +477,9 @@ Interactive notebooks (open in [Livebook](https://livebook.dev)):
- **[03 ADBC](livebook/03_adbc.livemd)** — Database, Connection, Statement, Stream (`:adbc_package` in Livebook).
- **[04 ADBC integration](livebook/04_adbc_integration.livemd)** — Connection pooling with NimblePool.
- **[05 Parquet](livebook/05_parquet.livemd)** — Pushdown reads, compressed writes, multi-file directories, PyArrow side-by-side.
- **[06 Datasets](livebook/06_datasets.livemd)** — Hive Dataset open, Expression scanner, prune stats, PyArrow `dataset` side-by-side.

See [livebook/README.md](livebook/README.md) for run instructions. Notebooks use Hex `~> 0.8.0` by default; opening from `livebook/` in a clone builds from source.
See [livebook/README.md](livebook/README.md) for run instructions. Notebooks use Hex `~> 0.9.0` by default; opening from `livebook/` in a clone builds from source.

---

Expand Down Expand Up @@ -763,6 +806,40 @@ Full guide: [docs/parquet_guide.md](docs/parquet_guide.md). Livebook:

---

## Datasets and scanners

Discover Hive-partitioned Parquet (or IPC) trees, then scan with projection
and Expression filters. Partition pruning, Parquet row-group pushdown, and
residual `Compute.filter/2` form a pushdown ladder.

```elixir
alias ExArrow.Compute.Expression, as: E

{:ok, dataset} =
ExArrow.Dataset.open("/data/events",
partitioning: {:hive, schema: [{"year", :int32}, {"month", :int32}]}
)

filter =
E.and_(
E.gte(E.field("year"), E.scalar(2026)),
E.gt(E.field("amount"), E.scalar(0.0))
)

{:ok, scanner} =
ExArrow.Dataset.scanner(dataset, columns: ["id", "amount"], filter: filter)

{:ok, stream} = ExArrow.Scanner.to_stream(scanner)
batches = Enum.to_list(stream)
ExArrow.Stream.close(stream)
ExArrow.Scanner.stats(stream)
```

Full guide: [guides/11_datasets.md](guides/11_datasets.md). Livebook:
`livebook/06_datasets.livemd`.

---

## Arrow compute kernels

All operations run entirely in native memory. Results are new
Expand Down Expand Up @@ -1082,6 +1159,7 @@ HTML reports are written to `bench/output/` (gitignored).
| `v070_stream_flow_pipeline_bench.exs` | Parquet/IPC stream drains, Flow execution, Pipeline map_batches + write_parquet at 1K/100K/1M rows |
| `v070_record_batch_vs_maps_bench.exs` | Arrow `RecordBatch` vs `list(map())` for build, transform, and drain at 1K/100K/1M rows |
| `parquet_pushdown_bench.exs` | Pushdown filter vs full read + project (v0.8.0 rough timing helper) |
| `dataset_scan_bench.exs` | Dataset pushdown ladder: full / project / partition / row-group / residual (v0.9.0) |


### Published results
Expand All @@ -1106,6 +1184,7 @@ The CI workflow posts a PR alert comment when any scenario regresses more than
- [Arrow and GenStage](guides/08_arrow_and_genstage.md) — demand-driven producers with backpressure
- [Arrow and Broadway](guides/09_arrow_and_broadway.md) — ingestion pipelines with `BatchBuilder` and sinks
- [Arrow Pipeline Patterns](guides/10_arrow_pipeline_patterns.md) — composable `ExArrow.Pipeline` transforms and sinks
- [Datasets and Scanners](guides/11_datasets.md) — Dataset, Fragment, Hive partitioning, Scanner pushdown ladder
- [Memory model](docs/memory_model.md) — handles, copying rules, NIF scheduling
- [IPC guide](docs/ipc_guide.md) — stream vs file, types, limitations
- [Parquet guide](docs/parquet_guide.md) — read/write Parquet, streaming, comparison with IPC
Expand Down Expand Up @@ -1303,6 +1382,18 @@ welcome for any of them.
- **Docs / CI** — Parquet Livebook, Mix.install smoke, Livebook pin checks;
optional arrow-testing corpus (local / workflow_dispatch).

### Shipped (v0.9.0)

- **Dataset / Fragment / Hive discovery** — `Dataset.open/2` over Local or
Memory filesystems; typed Hive partition values; schema from footer.
- **Scanner** — partition prune, Parquet pushdown, residual Expression filter;
`:dataset` stream backend; exact `Scanner.stats/1`.
- **`ExArrow.Compute.Expression`** — builders, schema validate, pushable vs
residual compile; `Compute.filter/2` residual NIF.
- **`RecordBatch.from_lists/1` / `from_map/1`** — list/map batch construction.
- **arrow-rs 59.3.0** — coordinated parquet / Flight / ADBC companion bumps.
- **Docs** — Datasets guide, Livebook 06, dataset scan bench.

### Longer-term

- **Streaming writes to Delta Lake** — sink for data pipeline nodes.
Expand Down
105 changes: 105 additions & 0 deletions bench/dataset_scan_bench.exs
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
# Dataset scan pushdown ladder — timing helper for announcements.
# Usage: mix run bench/dataset_scan_bench.exs
#
# Each label states exactly what the branch measures (F-013).
# Synthetic layout: >= 8 Hive partitions, 1M+ rows total.

alias ExArrow.Compute.Expression, as: E
alias ExArrow.Dataset
alias ExArrow.Parquet
alias ExArrow.RecordBatch
alias ExArrow.Scanner
alias ExArrow.Stream

partitions = 8
rows_per_part = 150_000
total_rows = partitions * rows_per_part

root = Path.join(System.tmp_dir!(), "ex_arrow_dataset_scan_bench")
File.rm_rf!(root)

IO.puts("Writing #{total_rows} rows across #{partitions} hive partitions under #{root}...")

Enum.each(0..(partitions - 1), fn p ->
year = 2020 + rem(p, 4)
month = rem(p, 12) + 1
n = rows_per_part

ids = for i <- 1..n, into: <<>>, do: <<i + p * n::little-signed-64>>
amounts = for i <- 1..n, into: <<>>, do: <<i * 1.0::little-float-64>>

{:ok, batch} =
RecordBatch.from_columns(["id", "amount"], [ids, amounts], ["s64", "f64"], n)

schema = RecordBatch.schema(batch)
path = Path.join(root, "year=#{year}/month=#{month}/part-0.parquet")
File.mkdir_p!(Path.dirname(path))
:ok = Parquet.Writer.to_file(path, schema, [batch], row_group_size: 50_000)
end)

{:ok, dataset} =
Dataset.open(root, partitioning: {:hive, schema: [{"year", :int32}, {"month", :int32}]})

measure = fn label, fun ->
{us, result} = :timer.tc(fun)
IO.puts("#{label}: #{Float.round(us / 1000, 1)} ms -> #{inspect(result)}")
end

measure.("full scan (all fragments, no filter, all columns)", fn ->
{:ok, scanner} = Dataset.scanner(dataset)
{:ok, stream} = Scanner.to_stream(scanner)
rows = Enum.sum(Enum.map(Enum.to_list(stream), &RecordBatch.num_rows/1))
Stream.close(stream)
rows
end)

measure.("projection only (columns: id — no filter)", fn ->
{:ok, scanner} = Dataset.scanner(dataset, columns: ["id"])
{:ok, stream} = Scanner.to_stream(scanner)
rows = Enum.sum(Enum.map(Enum.to_list(stream), &RecordBatch.num_rows/1))
Stream.close(stream)
rows
end)

year_filter = E.gte(E.field("year"), E.scalar(2022))

measure.("partition-pruned (year >= 2022; opens fewer fragments)", fn ->
{:ok, scanner} = Dataset.scanner(dataset, filter: year_filter)
{:ok, stream} = Scanner.to_stream(scanner)
rows = Enum.sum(Enum.map(Enum.to_list(stream), &RecordBatch.num_rows/1))
stats = Scanner.stats(stream)
Stream.close(stream)
{rows, stats.fragments_pruned_partition, stats.fragments_scanned}
end)

rg_filter =
E.and_(
E.gte(E.field("year"), E.scalar(2022)),
E.gt(E.field("amount"), E.scalar(140_000.0))
)

measure.("row-group-pruned (partition + amount > 140000 Parquet filters)", fn ->
{:ok, scanner} = Dataset.scanner(dataset, columns: ["id"], filter: rg_filter)
{:ok, stream} = Scanner.to_stream(scanner)
rows = Enum.sum(Enum.map(Enum.to_list(stream), &RecordBatch.num_rows/1))
stats = Scanner.stats(stream)
Stream.close(stream)
{rows, stats.row_groups_skipped, stats.row_groups_selected}
end)

residual_filter =
E.and_(
E.gte(E.field("year"), E.scalar(2022)),
E.not_(E.eq(E.field("id"), E.scalar(1)))
)

measure.("expression-residual (partition prune + not_ residual after decode)", fn ->
{:ok, scanner} = Dataset.scanner(dataset, columns: ["id"], filter: residual_filter)
{:ok, stream} = Scanner.to_stream(scanner)
rows = Enum.sum(Enum.map(Enum.to_list(stream), &RecordBatch.num_rows/1))
stats = Scanner.stats(stream)
Stream.close(stream)
{rows, stats.rows_emitted, stats.fragments_scanned}
end)

IO.puts("done.")
Loading