diff --git a/PROMPT.md b/PROMPT.md deleted file mode 100644 index d87a7ea..0000000 --- a/PROMPT.md +++ /dev/null @@ -1,210 +0,0 @@ - -You are a senior principal engineer and implementation agent. You will implement an Elixir library named **ExArrow** that provides complete Apache Arrow support for the BEAM, including: - -1) Arrow IPC (stream + file) -2) Arrow Flight (client + server) -3) ADBC (Arrow Database Connectivity) bindings + ergonomic Elixir API - -You must follow modern best practices in Elixir/OTP, Rust (arrow-rs, arrow-flight + tonic), and ADBC. You must ship Hex-quality code, tests, and docs. - -CRITICAL RULES -- DO NOT commit to git. Never run git commit, never assume commits. Instead, when a logical unit is complete, tell me: - - “Recommended commit message:” + a short imperative commit message - - A bullet list of what changed - - Wait for me to commit. -- Docs must be clean, elegant, and to the point. No emoji, no “decorative” writing, no marketing fluff. Use short sections, crisp explanations, minimal verbosity. -- Do not block BEAM schedulers from NIFs. Use dirty NIFs or native threads + message passing for long-running work. -- Keep Arrow data in Rust/Arrow buffers; Elixir should hold lightweight handles/resources. Avoid copying into the BEAM heap except when explicitly requested. -- Provide a stable, minimal core API first, then add ergonomic helpers. - -ADDITIONAL CONSTRAINTS (must enforce) -1) API Stability / Deprecation Policy - - No breaking API changes within a minor series (0.x: treat x as minor series; e.g., 0.2.* must not break). - - If a breaking change is unavoidable, implement a deprecation path: - - Keep old API for at least one minor series - - Emit warnings and document migration steps - - Provide explicit changelog entry and “Upgrade guide” notes - -2) Performance Gate (must be measurable and tested) - - Add a lightweight performance/heap-allocation gate to CI for critical paths: - - IPC stream encode/decode roundtrip for N rows must not allocate more than a specified threshold on the BEAM heap, and must complete within a reasonable time bound. - - Use a deterministic micro-benchmark-style ExUnit test (tagged, e.g. `@tag :perf`) that: - - Captures `:erlang.memory/0` (or `:erlang.memory(:total)`) before/after on the Elixir side - - Ensures large payloads stay mostly in native memory (handles only) - - Is conservative to avoid flakiness (skip/relax on slow CI if necessary but keep a baseline) - - Document exactly what is measured and why it correlates to “no big BEAM copies”. - -3) Test Driven Development (TDD) enforcement - - Use TDD by default: - - For each new feature, first add/modify ExUnit tests that define the behavior. - - Then implement the minimal code to pass the tests. - - Then refactor with tests still passing. - - In each milestone response, present work in this order: - 1) Tests (new/changed) with explanation of intended behavior - 2) Implementation code changes to satisfy tests - 3) Refactoring/cleanup - 4) Docs/examples updates - - Do not introduce large untested modules. Every public function must have direct or indirect test coverage. - - Use property tests (StreamData) for IPC roundtrip early, but keep generators bounded for stability. - -AUTHORITATIVE REFERENCES (use as design constraints) -- Arrow format + IPC spec: https://arrow.apache.org/docs/format/ -- Arrow Flight spec (gRPC + IPC): https://arrow.apache.org/docs/format/Flight.html -- ADBC overview + API standard, canonical in adbc.h: https://arrow.apache.org/adbc/ and https://arrow.apache.org/adbc/current/format/specification.html -- Rust Arrow crates (arrow-rs, arrow-ipc, arrow-flight): use current crates; keep versions aligned. - -PROJECT GOAL -Create ExArrow as foundational infrastructure, enabling interoperability with Python (PyArrow/Polars), data warehouses (Snowflake via Arrow), and DB connectivity (ADBC), with an OTP-friendly, production-grade API. - -DELIVERABLES (must be produced in your responses as you go) -- Repository tree and file purposes -- Public API outline (modules/functions) before implementing -- Detailed milestone plan with acceptance criteria (Definition of Done) -- For each milestone: code + tests + docs + examples -- A release checklist and CI matrix plan - -ENGINEERING WORKFLOW (branch-per-milestone delivery contract) -- Work in milestones. Each milestone has: - 1) Scope - 2) Acceptance criteria - 3) Implementation steps - 4) Tests - 5) Docs - 6) Example(s) - 7) Recommended commit message (but DO NOT commit) -- Keep changes incremental and reviewable. -- After each milestone, summarize status and prompt me to commit. - -REPO STRUCTURE (create exactly this shape) -- mix.exs, mix.lock -- lib/ex_arrow.ex (public entry, docs) -- lib/ex_arrow/{schema,field,array,record_batch,table,stream,error}.ex (as needed) -- lib/ex_arrow/ipc/*.ex -- lib/ex_arrow/flight/*.ex -- lib/ex_arrow/adbc/*.ex -- native/ex_arrow_native/Cargo.toml -- native/ex_arrow_native/src/*.rs -- test/* with unit + integration + property tests -- examples/ipc_roundtrip.exs -- examples/flight_echo/{server.exs,client.exs} -- examples/adbc_query.exs -- docs/ (guides) or ExDoc “Pages” configuration (preferred) - -CHOICE OF IMPLEMENTATION APPROACH -- Use Rustler NIFs for Arrow core: schema, arrays, record batches, IPC IO, Flight, and ADBC bindings. -- Use Resource types (opaque handles) on Elixir side: SchemaRef, ArrayRef, RecordBatchRef, StreamRef, FlightClientRef, AdbcDatabaseRef, etc. -- Return Elixir-friendly structures for small metadata, but keep large payloads native. - -QUALITY BAR -- Error handling: define ExArrow.Error with code/message/details; map Rust errors into structured Elixir errors. -- Telemetry: emit events for IPC read/write, Flight calls, ADBC execute, stream consumption, with durations and sizes. -- Types: support major Arrow types end-to-end: - Null, Bool, ints, floats, utf8/binary (large variants too), decimals, dates, timestamps (tz), lists, structs, dictionary, map (if feasible). -- Tests: - - Property tests for IPC roundtrip of bounded random schemas/batches - - Golden fixtures for IPC compatibility - - Flight integration tests (server + client) - - ADBC integration tests (at least one driver; tests must skip gracefully if driver not available) - -DOCS REQUIREMENTS (clean/elegant) -- ExDoc pages: - - Overview: what ExArrow is and what it is not - - IPC guide: stream vs file; sequential vs random access; examples - - Flight guide: client/server patterns; TLS; cancellation - - ADBC guide: concepts (Database/Connection/Statement), driver loading, examples - - Memory model: handles, lifetimes, copying rules - - Interop: how to exchange with Python via IPC (simple commands) -- Keep each guide short, factual, and actionable. No fluff. - -MILESTONES (you must follow these; do not skip) -Milestone 0: Skeleton + API surface -- Create the project structure, basic modules, and stubs. -- Define public API outline and resource-handle strategy. -- CI scaffold (at least Elixir + Rust build/test). -- DoD: - - `mix test` passes - - NIF compiles and loads (even if minimal) - - Docs build cleanly - -Milestone 1: IPC MVP (streaming only) -- Implement Schema + RecordBatch handles and IPC stream read/write: - - read from binary and from file path - - write to binary and to file path -- Types: bool, int64, float64, utf8, binary, nulls -- Streaming iterator in Elixir that yields RecordBatch handles -- DoD: - - Roundtrip tests (encode -> decode equals by schema + row count + sample values) - - Example `examples/ipc_roundtrip.exs` works - - IPC guide page exists and is concise - -Milestone 2: IPC Complete (file + more types) -- IPC file format (random access to batch i; schema read; batch count) -- Add nested types (list/struct), timestamps, decimals, dictionary encoding -- Golden fixtures strategy -- DoD: - - File random access tests - - Property tests added - - Documentation updated with explicit limitations (if any) - -Milestone 3: Flight MVP (client do_get/do_put) -- Flight client: - - connect (TLS optional) - - do_get -> stream batches - - do_put -> upload batches -- Minimal Flight server: - - Echo server supports do_put then do_get by ticket -- DoD: - - Integration test spins server, does put/get roundtrip - - Example `examples/flight_echo/*` runs - -Milestone 4: Flight Complete -- list_flights, get_flight_info, get_schema -- list_actions, do_action -- Timeouts, cancellation hooks, retry policy hooks (documented) -- DoD: - - Tests cover main Flight APIs - - Flight guide is concise and accurate - -Milestone 5: ADBC MVP -- Bind to ADBC driver manager (adbc.h canonical) and expose: - - Database.open - - Connection.open - - Statement.new(conn, sql) or new(conn, sql, bind: batch), execute; set_sql/bind for reuse/rebind - - execute returns an Arrow Stream (RecordBatch stream) -- Driver loading: - - load from shared library path and from env -- DoD: - - ADBC query example works against at least one driver (skip if unavailable) - - ADBC guide is short and practical - -Milestone 6: ADBC Complete -- Metadata APIs (tables/schemas/columns) where supported -- Parameter binding where supported -- Robust diagnostics + error mapping -- DoD: - - Extended tests (skip gracefully if driver not available) - - Docs updated with support matrix - -RELEASE CHECKLIST (must provide at end) -- Versioning -- CI matrix -- Rust/Elixir toolchain versions -- NIF build notes -- Compatibility notes (Arrow versions, ADBC driver manager) -- API stability/deprecation notes and upgrade guide pointer -- “Known limitations” list -- Performance gate description and how to run locally - -OUTPUT FORMAT FOR EACH RESPONSE YOU GIVE ME -For the current milestone: -1) Scope and acceptance criteria -2) Files you will create/modify (with paths) -3) Tests first (full code) + what they assert -4) Implementation code (full) -5) Refactors/cleanup notes -6) Docs page(s) updates (full) -7) How to run the example(s) -8) “Recommended commit message:” + bullet summary of changes (but do not commit) - -Begin now with Milestone 0. - diff --git a/README.md b/README.md index 0a04822..bb808df 100644 --- a/README.md +++ b/README.md @@ -31,6 +31,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.8.0](#whats-changed-in-v080) - [What's changed in v0.7.0](#whats-changed-in-v070) - [Livebook tutorials](#livebook-tutorials) - [IPC: stream and file](#ipc-stream-and-file) @@ -53,6 +54,7 @@ Native Apache Arrow for the BEAM: IPC streaming, Arrow Flight, Arrow Flight SQL, - [Shipped (v0.4.0)](#shipped-v040) - [Shipped (v0.5.0)](#shipped-v050) - [Shipped (v0.7.0)](#shipped-v070) + - [Shipped (v0.8.0)](#shipped-v080) - [FAQ](#faq) - [License](#license) @@ -106,11 +108,12 @@ Dirty NIF schedulers are used for blocking I/O. - **A uniform stream abstraction** — `ExArrow.Stream` works identically for IPC, Flight, and ADBC results. Code that processes batches does not know or care where the data came from. -- **Arrow-native pipelines (v0.7.0)** — `ExArrow.Pipeline` provides a lazy, +- **Arrow-native pipelines** — `ExArrow.Pipeline` provides a lazy, composable DSL for transforming and sinking streams of `RecordBatch` values. `ExArrow.Flow`, `ExArrow.GenStage`, and `ExArrow.Broadway` integrations cover parallel, demand-driven, and ingestion workloads. Telemetry events - fire at every stage. + fire at every stage. **Larger-than-memory Parquet** (pushdown, write options, + multi-file streams) lands in v0.8.0. --- @@ -250,8 +253,7 @@ and no Rust is needed: `Mix.install([{:ex_arrow, "~> 0.8.0"}])`. ## Quick start -**Read an Arrow stream with the v0.7.0 constructors** (one entry point for -every source): +**Read Parquet with pushdown** (v0.8.0 — one entry point for every source): ```elixir {:ok, stream} = ExArrow.Stream.from_parquet("/data/events.parquet") @@ -327,9 +329,43 @@ batch = ExArrow.Stream.next(stream) --- +## What's changed in v0.8.0 + +v0.8.0 focuses on **larger-than-memory Parquet**: read only the columns, row +groups, and rows you need, with peak memory scaling to the selected data rather +than the whole file. Pipeline modules from v0.7.0 remain the streaming layer on +top. + +### Parquet power-read / write + +| API | Purpose | +|-----|---------| +| `ExArrow.Stream.from_parquet/2` | Preferred entry: `:columns`, `:row_groups`, `:filters` pushdown | +| `ExArrow.Parquet.Reader` | Same opts on `from_file/2` / `from_binary/2`; `read_stats/1` for pruning | +| `ExArrow.Parquet.Writer` | `:compression` (`:snappy`, `:zstd`, `{:zstd, level}`, `:lz4`, `:gzip`, `:none`), `:row_group_size`, `:dictionary` | +| `ExArrow.Parquet.Metadata` | Footer-only metadata (row groups, column stats, encodings, kv) | +| `Stream.from_parquet_files/2`, `from_parquet_dir/2` | Lazy multi-file / directory streams; `Stream.close/1` for early abandon | +| `ExArrow.IPC.File.write/3`, `RecordBatch.concat/1` | Public wrappers for existing NIFs | + +Filter AST examples: `{:gt, "score", 0.9}`, +`{:and, [{:gte, "id", 1}, {:lt, "id", 1000}]}`. Row-group min/max statistics +prune when safe; remaining rows use parquet-rs `RowFilter`. + +### Docs, Livebook, CI + +- Parquet Livebook (`livebook/05_parquet.livemd`) and rewritten + [Parquet guide](https://ex-arrow.hexdocs.pm/parquet_guide.html) +- Pushdown timing helper: `bench/parquet_pushdown_bench.exs` +- Mix.install optional-dep smoke job; Livebook Hex-pin check; optional + apache/arrow-testing suite (local / `workflow_dispatch`) + +See [CHANGELOG.md](CHANGELOG.md) for the full list. + +--- + ## What's changed in v0.7.0 -v0.7.0 turns ExArrow from a transport and interchange library into the +v0.7.0 turned ExArrow from a transport and interchange library into the foundation for **Arrow-native data pipelines on the BEAM**. The central architectural principle: **operate on Arrow `RecordBatch` values** — not `list(map())`, not `Explorer.DataFrame`, not `Nx.Tensor`. Explorer and Nx @@ -387,11 +423,6 @@ without `:telemetry`, the Flow/GenStage/Broadway modules return - `bench/v070_record_batch_vs_maps_bench.exs` — Arrow `RecordBatch` vs `list(map())` for build, transform, and drain. -### Stats - -825 tests, 16 properties, 89.9% coverage. `mix format`, `mix credo --strict`, -`mix dialyzer`, `mix sobelow`, `mix docs --warnings-as-errors` all pass. - --- ## Livebook tutorials @@ -761,7 +792,7 @@ All operations run entirely in native memory. Results are new ## Batch operations -`ExArrow.Batch` (v0.7.0) provides lightweight `RecordBatch` transforms that +`ExArrow.Batch` provides lightweight `RecordBatch` transforms that stay in native Arrow memory. It is not a dataframe — use Explorer for analytics. All functions return `{:ok, batch} | {:error, message}`. @@ -788,7 +819,7 @@ analytics. All functions return `{:ok, batch} | {:error, message}`. ## Streaming pipelines -`ExArrow.Pipeline` (v0.7.0) is a thin, lazy abstraction for transforming and +`ExArrow.Pipeline` is a thin, lazy abstraction for transforming and sinking Arrow streams. Pipelines compose with `|>/2` and short-circuit on error. @@ -1048,8 +1079,9 @@ HTML reports are written to `bench/output/` (gitignored). | `pipeline_bench.exs` | End-to-end: IPC file on disk to Flight doput without materialising in BEAM | | `explorer_arrow_bench.exs` | Explorer <-> Arrow interchange at 1K/100K/1M rows | | `nx_arrow_bench.exs` | Nx <-> Arrow interchange at 1K/100K/1M rows (rank-1 and rank-2) | -| `v070_stream_flow_pipeline_bench.exs` | Parquet/IPC stream drains, Flow execution, Pipeline map_batches + write_parquet at 1K/100K/1M rows (v0.7.0) | -| `v070_record_batch_vs_maps_bench.exs` | Arrow `RecordBatch` vs `list(map())` for build, transform, and drain at 1K/100K/1M rows (v0.7.0) | +| `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) | ### Published results @@ -1068,8 +1100,8 @@ The CI workflow posts a PR alert comment when any scenario regresses more than - [Explorer Integration](guides/02_explorer_integration.md) — from_dataframe, to_dataframe, type mapping, limitations - [Nx Integration](guides/03_nx_integration.md) — from_nx, to_nx, boolean tensors, rank-2 - [Arrow Ecosystem](guides/04_arrow_ecosystem.md) — how ExArrow complements Explorer, Nx, ADBC, Flight, Parquet, ExZarr -- [Arrow Pipelines Overview](guides/05_arrow_pipelines_overview.md) — orientation to the v0.7.0 pipeline modules -- [Arrow Streams](guides/06_arrow_streams.md) — the v0.7.0 streaming abstraction, constructors, consumption patterns +- [Arrow Pipelines Overview](guides/05_arrow_pipelines_overview.md) — orientation to the pipeline modules +- [Arrow Streams](guides/06_arrow_streams.md) — the streaming abstraction, constructors, consumption patterns - [Arrow and Flow](guides/07_arrow_and_flow.md) — parallel batch processing with `ExArrow.Flow` - [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 @@ -1258,6 +1290,19 @@ welcome for any of them. - **Educational guides** — `guides/06..10` covering streams, Flow, GenStage, Broadway, and pipeline patterns. +### Shipped (v0.8.0) + +- **Parquet read pushdown** — `:columns`, `:row_groups`, `:filters` on + `Stream.from_parquet/2` and `Parquet.Reader`; statistics pruning + + `read_stats/1`. +- **Parquet write options** — `:compression`, `:row_group_size`, `:dictionary`. +- **`ExArrow.Parquet.Metadata`** — footer-only metadata (including encodings). +- **Multi-file streams** — `from_parquet_files/2`, `from_parquet_dir/2`, + `Stream.close/1`. +- **Public IPC/RecordBatch helpers** — `IPC.File.write/3`, `RecordBatch.concat/1`. +- **Docs / CI** — Parquet Livebook, Mix.install smoke, Livebook pin checks; + optional arrow-testing corpus (local / workflow_dispatch). + ### Longer-term - **Streaming writes to Delta Lake** — sink for data pipeline nodes. diff --git a/lib/ex_arrow/application.ex b/lib/ex_arrow/application.ex index 102f466..0f4c3d8 100644 --- a/lib/ex_arrow/application.ex +++ b/lib/ex_arrow/application.ex @@ -16,12 +16,19 @@ defmodule ExArrow.Application do end defp adbc_package_configured? do - adbc_module = Module.safe_concat(["Elixir", "Adbc", "Database"]) + case Application.get_env(:ex_arrow, :adbc_package) do + opts when is_list(opts) and opts != [] -> + adbc_database_loaded?() - Code.ensure_loaded?(adbc_module) && - case Application.get_env(:ex_arrow, :adbc_package) do - opts when is_list(opts) and opts != [] -> true - _ -> false - end + _ -> + false + end + end + + defp adbc_database_loaded? do + Code.ensure_loaded?(Module.safe_concat(["Elixir", "Adbc", "Database"])) + rescue + # Optional `:adbc` dep — atom may not exist yet; avoid Module.concat/2. + ArgumentError -> false end end diff --git a/livebook/03_adbc.livemd b/livebook/03_adbc.livemd index 0916372..1db1426 100644 --- a/livebook/03_adbc.livemd +++ b/livebook/03_adbc.livemd @@ -48,19 +48,19 @@ This notebook covers: opening a database, executing SQL, consuming the result st ### Documentation -- [ADBC guide](https://ex-arrow.hexdocs.pm/adbc_guide.html) — driver loading, `:adbc_package` backend, pooling, local testing -- API reference: [ExArrow.ADBC.Database](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html), [ExArrow.ADBC.Connection](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Connection.html), [ExArrow.ADBC.Statement](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html), [ExArrow.ADBC.DriverHelper](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.DriverHelper.html), [ExArrow.ADBC.ConnectionPool](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.ConnectionPool.html) -- Package overview: [ex-arrow.hexdocs.pm](https://ex-arrow.hexdocs.pm) +* [ADBC guide](https://ex-arrow.hexdocs.pm/adbc_guide.html) — driver loading, `:adbc_package` backend, pooling, local testing +* API reference: [ExArrow.ADBC.Database](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html), [ExArrow.ADBC.Connection](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Connection.html), [ExArrow.ADBC.Statement](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html), [ExArrow.ADBC.DriverHelper](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.DriverHelper.html), [ExArrow.ADBC.ConnectionPool](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.ConnectionPool.html) +* Package overview: [ex-arrow.hexdocs.pm](https://ex-arrow.hexdocs.pm) --- ## 1. Concepts -| Handle | Module | Purpose | -| ---------- | ---------------------------------------------------------------------- | ---------------------------------------------------------------- | -| Database | [`ExArrow.ADBC.Database`](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html) | Driver (path or name) + init options (e.g. URI). | +| Handle | Module | Purpose | +| ---------- | ------------------------------------------------------------------------------------- | ---------------------------------------------------------------- | +| Database | [`ExArrow.ADBC.Database`](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html) | Driver (path or name) + init options (e.g. URI). | | Connection | [`ExArrow.ADBC.Connection`](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Connection.html) | Session from a database. | -| Statement | [`ExArrow.ADBC.Statement`](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html) | Set SQL, execute → returns `ExArrow.Stream` of record batches. | +| Statement | [`ExArrow.ADBC.Statement`](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html) | Set SQL, execute → returns `ExArrow.Stream` of record batches. | Flow: **Database.open** → **Connection.open** → **Statement.new(conn, sql)** → **execute** → **Stream** (schema + next). diff --git a/livebook/04_adbc_integration.livemd b/livebook/04_adbc_integration.livemd index 93b33c7..6e69010 100644 --- a/livebook/04_adbc_integration.livemd +++ b/livebook/04_adbc_integration.livemd @@ -46,9 +46,9 @@ This notebook demonstrates the **adbc_package** backend: a single Livebook-frien ### Documentation -- [ADBC guide](https://ex-arrow.hexdocs.pm/adbc_guide.html) — `:adbc_package` backend and connection pooling -- API reference: [ExArrow.ADBC.Database](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html), [ExArrow.ADBC.Connection](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Connection.html), [ExArrow.ADBC.Statement](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html), [ExArrow.ADBC.ConnectionPool](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.ConnectionPool.html) -- Package overview: [ex-arrow.hexdocs.pm](https://ex-arrow.hexdocs.pm) +* [ADBC guide](https://ex-arrow.hexdocs.pm/adbc_guide.html) — `:adbc_package` backend and connection pooling +* API reference: [ExArrow.ADBC.Database](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Database.html), [ExArrow.ADBC.Connection](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Connection.html), [ExArrow.ADBC.Statement](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.Statement.html), [ExArrow.ADBC.ConnectionPool](https://ex-arrow.hexdocs.pm/ExArrow.ADBC.ConnectionPool.html) +* Package overview: [ex-arrow.hexdocs.pm](https://ex-arrow.hexdocs.pm) --- diff --git a/livebook/05_parquet.livemd b/livebook/05_parquet.livemd index c8d684e..eaeb153 100644 --- a/livebook/05_parquet.livemd +++ b/livebook/05_parquet.livemd @@ -1,29 +1,66 @@ # ExArrow — Parquet power-read / write (v0.8) ```elixir -Mix.install([ - {:ex_arrow, "~> 0.8.0"}, - {:kino, "~> 0.14"} -]) +deps = [ + {:pythonx, "~> 0.4.2"}, + {:kino_pythonx, "~> 0.1.0"}, + {:kino, "~> 0.19.0"} +] + +# Opened from livebook/ in the repo → local source; otherwise Hex (precompiled NIF). +local? = File.exists?(Path.join(__DIR__, "../native/ex_arrow_native/Cargo.toml")) + +{ex_arrow_dep, extra_deps, config} = + if local? do + System.put_env("EX_ARROW_BUILD", "1") + + # Force recompile ex_arrow so a cached Mix.install build is not reused. + ex_arrow_beam = Path.join(__DIR__, "../_build/dev/lib/ex_arrow/ebin") + + if File.dir?(ex_arrow_beam) do + File.rm_rf!(ex_arrow_beam) + end + + { + {:ex_arrow, path: Path.expand("..", __DIR__)}, + [{:rustler, "~> 0.36", optional: true}], + [rustler_precompiled: [force_build: [ex_arrow: true]]] + } + else + { + {:ex_arrow, "~> 0.8.0"}, + [], + [] + } + end + +Mix.install(deps ++ [ex_arrow_dep] ++ extra_deps, config: config) ``` -Docs: [Parquet guide](https://ex-arrow.hexdocs.pm/parquet_guide.html) · -[API](https://ex-arrow.hexdocs.pm) +```pyproject.toml +[project] +name = "project" +version = "0.0.0" +requires-python = "==3.13.*" +dependencies = ["pyarrow"] +``` -## PyArrow vs ExArrow — read a subset +## Overview -**PyArrow** +This notebook mirrors common PyArrow Parquet patterns with ExArrow v0.8: -```python -import pyarrow.parquet as pq -table = pq.read_table("events.parquet", columns=["id", "score"], - filters=[("score", ">", 0.9)]) -``` +* column projection + predicate pushdown +* compressed writes and footer metadata +* multi-file / directory streams +* IPC file-format write for symmetry -**ExArrow** +Run the cells **top to bottom** — later cells reuse `path`, `schema`, and `batch`. + +--- + +## Setup: write a small demo Parquet file ```elixir -# Build a small demo file {:ok, batch} = ExArrow.RecordBatch.from_columns( ["id", "score"], @@ -39,21 +76,55 @@ schema = ExArrow.RecordBatch.schema(batch) path = Path.join(System.tmp_dir!(), "ex_arrow_parquet_demo.parquet") :ok = ExArrow.Parquet.Writer.to_file(path, schema, [batch], compression: :zstd) +{path, ExArrow.RecordBatch.num_rows(batch), ExArrow.Schema.field_names(schema)} +``` + +--- + +## Read a subset (pushdown) + +**PyArrow equivalent:** + +```python +import pyarrow.parquet as pq + +# Elixir binaries arrive in Pythonx as `bytes` — decode before treating as a path. +parquet_path = path.decode("utf-8") if isinstance(path, (bytes, bytearray)) else path + +table = pq.read_table( + parquet_path, + columns=["id", "score"], + filters=[("score", ">", 0.9)], +) +table.num_rows +``` + +**ExArrow:** + +```elixir {:ok, stream} = ExArrow.Stream.from_parquet(path, columns: ["id", "score"], filters: {:gt, "score", 0.9} ) -ExArrow.Parquet.Reader.read_stats(stream) -|> IO.inspect(label: "pushdown stats") +stats = ExArrow.Parquet.Reader.read_stats(stream) +rows = + stream + |> ExArrow.Stream.to_list() + |> Enum.map(&ExArrow.RecordBatch.num_rows/1) + |> Enum.sum() -ExArrow.Stream.to_list(stream) -|> Enum.map(&ExArrow.RecordBatch.num_rows/1) -|> IO.inspect(label: "filtered batch rows") +%{stats: stats, filtered_rows: rows} ``` -## Compressed write +Both cells return `2` — the rows where `score > 0.9` (`0.95` and `1.2`). Livebook +shares the Elixir `path` binding into the Python cell automatically; Pythonx encodes +Elixir binaries as Python `bytes`, so decode to `str` before passing to PyArrow. + +--- + +## Compressed write + footer metadata ```elixir :ok = @@ -63,14 +134,23 @@ ExArrow.Stream.to_list(stream) ) {:ok, meta} = ExArrow.Parquet.Metadata.from_file(path) -IO.inspect(meta.num_row_groups, label: "row groups") -IO.inspect(hd(meta.row_groups).columns, label: "column stats") + +%{ + num_rows: meta.num_rows, + num_row_groups: meta.num_row_groups, + columns: hd(meta.row_groups).columns +} ``` +Column maps include `path`, `compression`, `encodings`, `min`, and `max`. + +--- + ## Partitioned directory (multi-file) ```elixir dir = Path.join(System.tmp_dir!(), "ex_arrow_parquet_parts") +File.rm_rf!(dir) File.mkdir_p!(dir) for {name, id} <- [{"part-0.parquet", 10}, {"part-1.parquet", 20}] do @@ -91,14 +171,31 @@ for {name, id} <- [{"part-0.parquet", 10}, {"part-1.parquet", 20}] do end {:ok, multi} = ExArrow.Stream.from_parquet_dir(dir, filters: {:gte, "id", 15}) -IO.inspect(Enum.map(ExArrow.Stream.to_list(multi), &ExArrow.RecordBatch.num_rows/1)) +row_counts = Enum.map(ExArrow.Stream.to_list(multi), &ExArrow.RecordBatch.num_rows/1) + +# Early abandon releases the multi-file Agent. +:ok = ExArrow.Stream.close(multi) + +%{opened_dir: dir, filtered_batch_rows: row_counts} ``` +Expect one batch with `1` row (`id == 20`). + +--- + ## IPC file writer symmetry ```elixir -ipc_path = Path.join(System.tmp_dir!(), "demo.arrow") +ipc_path = Path.join(System.tmp_dir!(), "ex_arrow_parquet_demo.arrow") :ok = ExArrow.IPC.File.write(ipc_path, schema, [batch]) {:ok, file} = ExArrow.IPC.File.from_file(ipc_path) -ExArrow.IPC.File.batch_count(file) + +%{ + path: ipc_path, + batch_count: ExArrow.IPC.File.batch_count(file), + fields: + case ExArrow.IPC.File.schema(file) do + {:ok, sch} -> ExArrow.Schema.field_names(sch) + end +} ```