I/O Guide

Files, arguments, databases, and streams

Ibex keeps file I/O in plugins, not in core syntax. That means CSV and Parquet read/write functions, filesystem enumeration, command-line parsing, SQLite reads through ADBC, and Kafka streams are loaded with import and then called like any other function in scripts, the REPL, notebooks, and transpiled binaries.

What ships today

CSV plugin

Read RFC 4180 CSV into a DataFrame, infer types, accept custom null tokens, custom delimiters, headerless files, and optional schema hints when you already know the column types.

First-party Parquet backend

Read and write Apache Parquet for typed batch storage. This is the right default when you want a compact columnar format, better type preservation, and faster reloads than CSV. The standard hosts link this backend directly; import "parquet" does not require loading a shared plugin.

ADBC plugin

Read from SQLite and other ADBC-backed systems as Arrow streams. This is the bridge from database result sets into Ibex chunked execution without writing a bespoke reader per backend.

Kafka plugin

Consume JSON or Schema-Registry-backed Avro messages from Kafka-compatible brokers such as Redpanda, turn them into one-row typed tables, and plug them directly into Stream { ... } pipelines or WebSocket dashboards.

Filesystem plugin

fs::list returns directory entries in path order. Its fixed schema includes path, name, stem, extension, size, and directory status, making it a natural input to row-wise map jobs.

Arguments plugin

args::parse validates a declarative option specification and returns one row per option, flag, or positional. Optional missing values pair with nullable scalar extraction instead of sentinel strings.

Browse the filesystem, arguments, CSV, JSON, Parquet, ADBC, and Kafka plugin references.

Import bundled I/O backends

The common pattern is to import the bundled backend modules you need and then call their functions directly.

import "csv";
import "parquet";
import "fs";
import "args";
import "adbc";
import "kafka";

let trades = csv::read("data/trades.csv");
let prices = parquet::read("data/prices.parquet");
let public_prices = parquet::read("https://data.example.com/prices.parquet");
let cloud_prices = parquet::read("s3://market-data/prices.parquet?region=us-east-1");
let db_prices = adbc::read("/path/to/libadbc_driver_sqlite.so",
                               "data/prices.sqlite",
                               "select * from prices");
let tick = kafka::recv("localhost:19092", "ticks",
                          "ibex-demo",
                          "ts:timestamp,symbol:str,price:f64,size:i64");

For Schema Registry-backed Avro, use kafka::recv_avro(..., schema, registry_url[, options]).

Turn directories and command lines into ordinary tables

Convert every matching file

Row-wise map processes entries in input order. Each field is scalar-valued, but an effectful field may compose a reader with a writer. The result below records one row per source file.

import "csv";
import "fs";
import "parquet";

fs::list("incoming", "*.csv")[map {
    source = path,
    rows = parquet::write(
        csv::read(path),
        `archive/${stem}.parquet`
    )
}];

Parse script arguments

args::parse always returns kind, name, index, value; value is a string. A missing optional argument emits no row, so the one-argument scalar form yields null and coalesce can apply a script-level default. Pass script arguments after --.

import "args";

let args = args::parse("
    limit (n) : int?
    input     : positional+
");
let raw = scalar(args[filter name == "limit", select { value }]);
let limit = Int64(coalesce(raw, "100"));

// Run as: ibex job.ibex -- --limit 25 input.parquet

Read CSV with as much structure as you know

Default CSV read

With just a path, Ibex reads the header row, infers Int64, Float64, or string/categorical columns, and returns a DataFrame.

import "csv";

let trades = csv::read("data/trades.csv");

Null tokens

Pass a comma-separated null specification. <empty> means empty fields should be treated as nulls too.

import "csv";

let train = csv::read("train.csv", "<empty>,NA");

Custom delimiter

This is useful for semicolon-separated files or other non-comma dialects. Quoted commas are still preserved correctly.

import "csv";

let raw = csv::read("quotes.txt", "", ";");

Headerless files

When has_header is false, Ibex synthesizes col1, col2, and so on. This is what the 1BRC benchmark file uses.

import "csv";

let raw = csv::read("measurements.txt", "", ";", false)
    [select { station = col1, temp = col2 }];

Schema hints

If you already know the CSV schema, pass it explicitly to skip most of the inference work. For large repeated loads, this is the right shape.

import "csv";

// positional schema
let fast = csv::read("measurements.txt", "", ";", false, "cat,f64");

// named schema
let typed = csv::read("trades.csv", "", ",", true,
                          "symbol:cat,price:f64,size:i64");

// exact decimals: parsed straight to Decimal(12, 2), never through a double;
// extra fractional digits round half away from zero, too many digits is an error
let ledger = csv::read("ledger.csv", "", ",", true,
                           "account:str,amount:decimal(12,2)");

Run visible live-stream demos with Kafka input and WebSocket output

JSON demo

The bundled Kafka demo uses Redpanda as a local Kafka-compatible broker. A synthetic tick producer writes JSON ticks into the ticks topic, and two Ibex stream jobs consume that same topic:

  • one job groups by symbol and venue and sends live summaries to ws://127.0.0.1:8765
  • one job resamples ticks into 5-second OHLC bars and sends them to ws://127.0.0.1:8766

The dashboard at demo/kafka/ws_dashboard.html opens both WebSocket feeds and shows them as two tabs.

Start the JSON dashboard

# shell
demo/kafka/demo-kafka.sh
demo/kafka/run-kafka-dashboard.sh

# optional terminal viewers
python3 demo/kafka/ws_client.py
python3 demo/kafka/ohlc_ws_client.py

The Docker stack starts Redpanda plus the random tick producer. The launcher script starts two Ibex REPL processes in the background and keeps them running until you stop it.

The JSON Ibex side

import "kafka";
import "websocket";

ws::listen(8765);

let summary = Stream {
    source = kafka::recv(
        brokers = "localhost:19092",
        topic = "ticks",
        group = "ibex-demo-summary",
        schema = "ts:timestamp,symbol:str,venue:str,price:f64,size:i64",
        options = "poll_timeout_ms=100;consumer.auto.offset.reset=latest;consumer.session.timeout.ms=6000"
    ),
    transform = [by { symbol, venue }, select {
        trades = count(),
        last_px = last(price),
        total_size = sum(size)
    }],
    sink = ws::send(8765)
};

The OHLC job in demo/kafka/kafka_ohlc.ibex uses the same source but switches the transform to [resample 5s, by { symbol }, select { ... }] and writes to ws://127.0.0.1:8766.

Avro demo with Schema Registry

The same Docker stack also starts Redpanda Schema Registry on http://127.0.0.1:18081 and a second producer that writes Avro ticks into ticks_avro. The Avro dashboard uses two separate websocket feeds:

  • ws://127.0.0.1:8775 for grouped trade summaries
  • ws://127.0.0.1:8776 for 5-second OHLC bars
# shell
demo/kafka/demo-kafka.sh
demo/kafka/run-kafka-avro-dashboard.sh

# optional terminal viewers
python3 demo/kafka/ws_client.py --port 8775
python3 demo/kafka/ohlc_ws_client.py --port 8776

Open demo/kafka/ws_dashboard_avro.html to view the Avro dashboard in the browser.

The Avro Ibex side

import "kafka";
import "websocket";

ws::listen(8775);

let summary = Stream {
    source = kafka::recv_avro(
        brokers = "localhost:19092",
        topic = "ticks_avro",
        group = "ibex-demo-summary-avro",
        schema = "ts:timestamp,symbol:str,venue:str,price:f64,size:i64",
        registry_url = "http://localhost:18081",
        options = "poll_timeout_ms=100;consumer.auto.offset.reset=latest;consumer.session.timeout.ms=6000"
    ),
    transform = [by { symbol, venue }, select {
        trades = count(),
        last_px = last(price),
        total_size = sum(size)
    }],
    sink = ws::send(8775)
};

The matching OHLC job in demo/kafka/kafka_ohlc_avro.ibex reads the same ticks_avro topic and writes to ws://127.0.0.1:8776.

Kafka-specific notes

  • The schema string is explicit and required: use name:type entries like "ts:timestamp,symbol:str,venue:str,price:f64,size:i64".
  • For live dashboards, prefer consumer.auto.offset.reset=latest so the stream starts at current traffic instead of replaying old messages.
  • poll_timeout_ms is only the wait between polls. It is not a lifetime timeout for the stream.
  • For per-message live streams, str is a better demo key type than cat; it avoids per-message categorical dictionary churn.
  • kafka::recv_avro(...) requires a Schema Registry base URL and currently supports flat records only.
  • The Avro v1 path does not support schema references or null-valued unions yet.

Write CSV back out

csv::write writes a header row plus all data rows and returns the number of data rows written. Nulls become empty fields and strings are quoted according to RFC 4180.

import "csv";

let summary = trades[
    select { avg_px = mean(price), rows = count() },
    by symbol,
    order symbol
];

let n = csv::write(summary, "out/summary.csv");

Use Parquet for typed, columnar storage

Read Parquet

Parquet is the cleaner choice when you control the data pipeline and want strong type preservation, compact files, faster reloads than CSV, and the same reader for disk, public HTTPS URLs, presigned URLs, or S3-compatible object storage.

Ibex reads referenced Parquet columns directly into Ibex storage. For a row-local filter over a single scan, it decodes predicate columns first and late-materializes the other columns only for surviving row indices. Reusing the same binding in multiple scan positions safely falls back to ordinary column projection.

import "parquet";

let prices = parquet::read("data/prices.parquet");
let liquid = prices[filter size > 1000, select { ts, symbol, price }];
let public_prices = parquet::read("https://data.example.com/prices.parquet");
let cloud = parquet::read("s3://market-data/prices.parquet?region=us-east-1");

Write Parquet

parquet::write returns the number of rows written. String and categorical columns are stored as UTF-8. Dates and timestamps are preserved as typed Parquet columns, and Decimal(p, s) columns are written as Parquet DECIMAL with their precision and scale, so the values read back exactly. Reading accepts every physical encoding of a Parquet decimal up to 38 digits; a wider (Decimal256) column is refused by name rather than approximated.

import "parquet";

let n = parquet::write(prices, "out/prices.parquet");

Read database tables and query results through the chunked path

Minimal SQLite query

adbc::read takes a driver, a database URI or path, and a SQL query. The driver is a name such as "sqlite", resolved through an ADBC driver manifest, or a path to the driver library. scripts/install_adbc_driver.sh sqlite installs Apache's SQLite driver (pinned and SHA-256 checked, no Python or conda) and writes the manifest; postgresql, duckdb and mysql install the same way, the last two from the public driver registry dbc uses, without needing dbc. On Windows, run powershell -ExecutionPolicy Bypass -File scripts\install_adbc_driver.ps1 sqlite; the bypass lasts for that one process and is needed because Windows refuses unsigned scripts that came from a download.

Columns arrive as the matching Ibex type; float32 widens to Float64 and a uuid becomes its canonical text. A column Ibex has no type for (time, intervals, binary, arrays) is refused with an error that names it and the cast that reads it, such as CAST(t AS TEXT) (AS CHAR on MySQL, which has no TEXT cast). PostgreSQL's driver hands numeric over as text; convert it with Decimal(x, p, s).

When using import "adbc" in the REPL, start Ibex with --plugin-path ./build-release/tools so it can find both adbc.ibex and adbc.so.

import "adbc";

let trades = adbc::read(
    "sqlite",
    "/tmp/trades.sqlite",
    "select symbol, qty, px from trades order by qty desc"
);

Reuse one connection

adbc::read opens and closes a connection on every call. adbc::connect opens one that later statements share: one login, and session state such as temporary tables. adbc::query runs a query on it and returns a table, which can be used as any other table. adbc::close closes the connection for every binding of it, returning 1 the first time and 0 after; a query on a closed connection is an error. A connection also closes when its last binding goes away.

Connection functions run one statement at a time, in order, not inside a query clause. Functions can take, open and return connections. A query's result is read in full before the statement uses it. The connections guide covers the rules in full.

import "adbc";

let db = adbc::connect("sqlite", ":memory:");
adbc::execute(db, "create table trades (symbol text, qty integer)");
adbc::execute(db, "insert into trades values ('AAPL', 10), ('MSFT', 5)");   // 2

adbc::query(db, "select symbol, qty from trades")[
    select { total = sum(qty) }, by symbol
];

adbc::close(db);

adbc::connect(driver, uri, options = "") takes the same options string as adbc::read; stmt. options apply to every query on the connection. examples/adbc_connection.ibex is a runnable version.

Write tables and run statements

adbc::execute(db, sql) runs a statement that returns no rows (DDL, insert, update, delete) and returns the affected-row count, or -1 when the driver reports none. For DDL the count is the driver's: PostgreSQL reports -1, SQLite repeats the previous statement's count.

adbc::write(db, df, table, mode = "create") bulk-inserts a DataFrame through the driver's ingestion path (PostgreSQL uses COPY) and returns the rows written. mode is create (the table must not exist), append (it must), replace (drop and recreate) or create_append. A write that fails, say on a duplicate key, writes no rows; DuckDB's driver drops the reason, so there the error only says that the write failed. MySQL's driver inserts in batches of 1000, so there Ibex runs the rows in a transaction of their own; as MySQL commits whenever it creates or drops a table, a failed create or replace leaves the new table, empty. For the same reason, inside adbc::begin MySQL takes append only. DuckDB and MySQL append whole rows only: df needs every column of an existing table.

import "adbc";
import "csv";

let db = adbc::connect("postgresql", "postgresql://localhost/markets");
let trades = csv::read("data/trades.csv");
adbc::write(db, trades, "trades", "replace");
adbc::execute(db, "create index on trades (symbol)");

// Values bind to placeholders ($1 for PostgreSQL, ? for SQLite, MySQL): one
// prepared statement, one run per row, no quoting in the SQL text.
adbc::execute(db, "delete from trades where symbol = $1",
             Table { symbol = ["TSLA", "GME"] });
adbc::query(db, "select * from trades where qty > $1", Table { min_qty = [100] });
adbc::close(db);

adbc::query and adbc::execute take an optional parameter table. Its columns bind by position; each row runs the prepared statement once, a query's results following one another and adbc::execute's counts adding up. A table with no rows runs nothing.

adbc::tables(db) lists the tables and views a connection sees, and adbc::table_schema(db, table) gives each column's Arrow type and the Ibex type a query would give it, or why it has none, before you run a query; see Discovery.

adbc::begin(db), adbc::commit(db) and adbc::rollback(db) make the calls between them one transaction. If any of them fails, adbc::commit rolls back instead, on every driver, and closing a connection never commits; see Transactions.

Every Ibex column type has a PostgreSQL column type: Int64 as bigint, Float64 as double precision, Bool, String and Categorical as boolean and text, Date, Timestamp as timestamp (microseconds, PostgreSQL's precision), and Decimal as numeric. SQLite has no bool, date or timestamp type: its driver stores 0/1 and ISO 8601 text, and it refuses Decimal columns; convert them with Float64(x) first if that loss is acceptable.

Runnable trades walkthrough

examples/adbc_sqlite/ is self-contained: it builds an 11-row trades database with nulls, reads it in three Arrow batches, filters, computes notional and VWAP per symbol in Ibex, exports a CSV, and shows that an empty result keeps its schema. One command runs it and checks the output; it is also the adbc:sqlite_demo test.

scripts/install_adbc_driver.sh sqlite
examples/adbc_sqlite/run.sh build

1M-row measurements showcase

The helper script below loads the first 1,000,000 rows of the 1BRC benchmark file into SQLite. After that, the Ibex query shape is the same aggregation you would run directly over CSV.

# shell
python3 scripts/import_measurements_sqlite.py \
    --input examples/measurements.txt \
    --limit 1000000 \
    --output /tmp/measurements_1m.sqlite

LD_LIBRARY_PATH=~/envs/ibex/lib:$LD_LIBRARY_PATH \
./build-release/tools/ibex --plugin-path ./build-release/tools

# ibex
import "adbc";

let measurements = adbc::read(
    "/path/to/libadbc_driver_sqlite.so",
    "/tmp/measurements_1m.sqlite",
    "select station, temp from measurements"
);

measurements[select {
    min_temp = min(temp),
    avg_temp = mean(temp),
    max_temp = max(temp)
}, by station, order station];

Optional ADBC options string

The 4th argument is optional (it defaults to ""): a key=value string separated by semicolons or newlines. Prefix keys with db., conn., or stmt. to target database, connection, or statement options. The rest of the key reaches the driver unchanged, so use the full ADBC name: conn.adbc.connection.autocommit, not conn.autocommit.

conn.post. sets a connection option after the connection is opened, for options a driver only accepts on an open connection (SQLite extension loading, for example). A backslash escapes ;, = and \ inside keys and values. Option names beyond the standard adbc.* keys are driver-specific; check the driver's documentation.

let trades = adbc::read(
    "sqlite",
    "/tmp/trades.sqlite",
    "select symbol, qty, px from trades",
    "stmt.adbc.sqlite.query.batch_rows=10000;conn.post.adbc.sqlite.load_extension.enabled=true"
);

Typical batch pipeline

import "csv";
import "parquet";

let raw = csv::read("measurements.txt", "", ";", false, "cat,f64")
    [select { station = col1, temp = col2 }];

let summary = raw[select {
    min_temp = min(temp),
    avg_temp = mean(temp),
    max_temp = max(temp)
}, by station, order station];

parquet::write(summary, "out/summary.parquet");

When to choose which format

Choose CSV when

  • You are ingesting raw files from outside systems.
  • You need human-readable interchange.
  • You need delimiter or headerless-file flexibility.

Choose Parquet when

  • You control both ends of the pipeline.
  • You want better typed persistence and faster reloads.
  • You are caching a cleaned dataset for repeated analysis.

Choose SQLite via ADBC when

  • You want to stage a database table or query result into Ibex.
  • You want SQL pushdown for filtering, joins, or pre-aggregation.
  • You want to bridge into other ADBC-backed systems later.

Choose Kafka when

  • You want live event ingestion instead of batch file loads.
  • You want to feed Stream { ... } pipelines from a broker rather than UDP or WebSocket input.
  • You want a local Redpanda-backed demo that ends in a browser dashboard.