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.
I/O Guide
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.
Overview
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.
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.
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.
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.
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.
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 pattern
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]).
Batch orchestration
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`
)
}];
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
CSV Input
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");
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");
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", "", ";");
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 }];
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)");
Kafka via Redpanda
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:
symbol and venue and sends live summaries to ws://127.0.0.1:8765ws://127.0.0.1:8766
The dashboard at demo/kafka/ws_dashboard.html opens both
WebSocket feeds and shows them as two tabs.
# 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.
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.
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 summariesws://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.
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.
name:type entries like "ts:timestamp,symbol:str,venue:str,price:f64,size:i64".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.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.CSV Output
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");
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");
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");
SQLite via ADBC
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"
);
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.
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.
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
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];
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"
);
Round Trip
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");
Practical Notes
Stream { ... } pipelines from a broker rather than UDP or WebSocket input.