Ingestion patterns: five sources, one table

Real desks never get data from one place: the tick feed hands you Arrow batches, research notebooks live in pandas or polars, vendors drop Parquet, and the odd legacy process still emails CSV. h5i-db's ingestion surface is deliberately small, append (extend the feed) and write (replace the contents), and both accept anything Arrow-shaped. This recipe ingests five consecutive trading days from five different sources into one trades table, then covers the operational half of ingestion: write vs append semantics, optimistic locking with expected_version, and why you batch commits and compact afterwards.

import shutil
from pathlib import Path

import pandas as pd
import polars as pl
import pyarrow as pa
import pyarrow.csv as pacsv
import pyarrow.parquet as pq

import h5i_db
import cookbook_utils as cu

db = h5i_db.Database(cu.fresh_db("00_ingestion"), create=True)

SCHEMA = pa.schema(
    [
        pa.field("ts", pa.timestamp("us", tz="UTC"), nullable=False),
        pa.field("symbol", pa.string()),
        pa.field("price", pa.float64()),
        pa.field("size", pa.int64()),
        pa.field("exchange", pa.string()),
        pa.field("side", pa.string()),
    ]
)
db.create_table("trades", SCHEMA, time_column="ts", sort_key=["ts", "symbol"])
output
{'table': 'trades',
 'sequence': 0,
 'op': 'create',
 'rows_total': 0,
 'segments_total': 0,
 'segments_added': 0,
 'segments_deduped': 0,
 'committed_at_ns': 1784778925550961624}

One continuous 11-session tape, sliced into per-day batches so each "delivery" arrives in time order, since append requires every batch to start at or after the table's stored max timestamp (feed semantics). A staging directory under data/dbs stands in for the vendor drop zone.

tape = cu.make_trades(days=11, trades_per_day=4_000, start="2026-06-01", seed=7)

dates = tape["ts"].to_pandas().dt.date
sessions = sorted(dates.unique())
by_day = {d: tape.filter(pa.array((dates == d).to_numpy())) for d in sessions}
print(f"{len(sessions)} sessions, {len(tape):,} trades:", sessions[0], "→", sessions[-1])

staging = Path("data/dbs/00_ingestion_staging")
if staging.exists():
    shutil.rmtree(staging)
staging.mkdir(parents=True)
output
11 sessions, 141,301 trades: 2026-06-01 → 2026-06-15

1. From a pyarrow Table#

The native path, with zero conversion. Every append returns the commit dict: the new version number (sequence), total rows and segment counts after the commit. Log it in your loaders; it is the receipt that ties a delivery to a version.

commit = db.append("trades", by_day[sessions[0]], note="day 1: arrow feed")
commit
output
{'table': 'trades',
 'sequence': 1,
 'op': 'append',
 'rows_total': 11418,
 'segments_total': 1,
 'segments_added': 1,
 'segments_deduped': 0,
 'committed_at_ns': 1784778925721745066}

2. From a pandas DataFrame#

pa.Table.from_pandas does the heavy lifting, but always pass schema=. Depending on your pandas version, datetimes round-trip as nanoseconds (not the table's microseconds) and strict append would refuse the mismatch. Passing the target schema makes the conversion do the cast.

df = by_day[sessions[1]].to_pandas()  # pretend this came from research code
commit = db.append(
    "trades",
    pa.Table.from_pandas(df, schema=SCHEMA, preserve_index=False),
    note="day 2: pandas",
)
{k: commit[k] for k in ("sequence", "rows_total", "segments_total")}
output
{'sequence': 2, 'rows_total': 25773, 'segments_total': 2}

3. From a polars DataFrame#

Polars speaks Arrow natively, with one wrinkle: its to_arrow() emits large_string columns, which strict append rejects as a schema mismatch. .cast(SCHEMA) is a cheap metadata-level fix; make it a habit on any polars → h5i-db boundary.

pldf = pl.from_arrow(by_day[sessions[2]])  # pretend this came from a polars pipeline
commit = db.append("trades", pldf.to_arrow().cast(SCHEMA), note="day 3: polars")
{k: commit[k] for k in ("sequence", "rows_total", "segments_total")}
output
{'sequence': 3, 'rows_total': 38842, 'segments_total': 3}

4. From a Parquet file#

Parquet preserves types exactly, so a vendor Parquet drop is a read-and-append. (If a vendor's row order is suspect, sort before appending; strict append will tell you either way.)

pq.write_table(by_day[sessions[3]], staging / "vendor_day4.parquet")

commit = db.append("trades", pq.read_table(staging / "vendor_day4.parquet"), note="day 4: parquet drop")
{k: commit[k] for k in ("sequence", "rows_total", "segments_total")}
output
{'sequence': 4, 'rows_total': 49876, 'segments_total': 4}

5. From CSV#

CSV keeps values but drops types: read it naively and ts comes back as whatever the parser guesses. ConvertOptions(column_types=...) pins the time column to timestamp[us, tz=UTC] at parse time, giving you back an exact schema match.

pacsv.write_csv(by_day[sessions[4]], staging / "legacy_day5.csv")

from_csv = pacsv.read_csv(
    staging / "legacy_day5.csv",
    convert_options=pacsv.ConvertOptions(column_types={"ts": pa.timestamp("us", tz="UTC")}),
)
assert from_csv.schema.equals(by_day[sessions[4]].schema)
commit = db.append("trades", from_csv, note="day 5: legacy csv")
{k: commit[k] for k in ("sequence", "rows_total", "segments_total")}
output
{'sequence': 5, 'rows_total': 63593, 'segments_total': 5}

write vs append#

Neither is ever destructive; old versions remain readable via db.read(..., version=n) and h5i('table', n) in SQL.

universe_schema = pa.schema(
    [pa.field("ts", pa.timestamp("us", tz="UTC"), nullable=False), pa.field("symbol", pa.string())]
)
db.create_table("universe", universe_schema, time_column="ts")

asof = pa.scalar(pd.Timestamp("2026-06-01", tz="UTC"), type=pa.timestamp("us", tz="UTC"))
db.write(
    "universe",
    pa.table({"ts": pa.array([asof] * 3), "symbol": pa.array(["AAPL", "MSFT", "NVDA"])}),
    note="June universe",
)
db.write(
    "universe",
    pa.table({"ts": pa.array([asof] * 4), "symbol": pa.array(["AAPL", "MSFT", "NVDA", "AVGO"])}),
    note="June universe, AVGO added",
)
print("head :", db.read("universe")["symbol"].to_pylist())
print("v1   :", db.read("universe", version=1)["symbol"].to_pylist())
[{k: v[k] for k in ("sequence", "op", "rows", "note") if k in v} for v in db.versions("universe")]
output
head : ['AAPL', 'MSFT', 'NVDA', 'AVGO']
v1   : ['AAPL', 'MSFT', 'NVDA']
[{'sequence': 0, 'op': 'create', 'rows': 0},
 {'sequence': 1, 'op': 'write', 'rows': 3, 'note': 'June universe'},
 {'sequence': 2,
  'op': 'write',
  'rows': 4,
  'note': 'June universe, AVGO added'}]

Optimistic locking with expected_version#

When two loaders share a table, "append whatever, whenever" silently interleaves deliveries. append(..., expected_version=n) is compare-and-swap: the commit only lands if the table head is still at version n, otherwise you get a ConflictError. Note retryable=True and the hint spelling out the recovery. The retry is mechanical: re-read the head, re-append.

day6 = by_day[sessions[5]]
try:
    db.append("trades", day6, expected_version=1)  # stale: head is already at v5
except h5i_db.ConflictError as e:
    print(f"code      {e.code}")
    print(f"retryable {e.retryable}")
    print(f"hint      {e.hint}")

# retry pattern: re-read the head version, then re-append against it
head = db.versions("trades")[-1]["sequence"]
commit = db.append("trades", day6, expected_version=head, note="day 6: CAS append")
print(f"\nretried against v{head} -> committed v{commit['sequence']}")
output
code      version_conflict
retryable True
hint      re-read the head of "trades" and retry against it; pure appends rebase safely (the CLI and Python bindings already auto-retry those)

retried against v5 -> committed v6

Batching and compaction#

Every commit writes a manifest and at least one segment, so commit batches (a day, an hour, a few thousand rows), never per row. But even day-sized commits accumulate small segments, and query planning touches every one of them. The daily-loop-then-compact pattern below is the normal rhythm: many small append commits during the week, one compact to merge segments. Compaction is itself just another commit: same data, fewer segments, full history retained.

for d in sessions[6:]:
    commit = db.append("trades", by_day[d], note=f"daily load {d}")
    print(f"v{commit['sequence']}: +{len(by_day[d]):>6,} rows "
          f"-> {commit['segments_total']:>2} segments total")
output
v7: +13,806 rows ->  7 segments total
v8: +12,169 rows ->  8 segments total
v9: +11,128 rows ->  9 segments total
v10: +13,093 rows -> 10 segments total
v11: +14,464 rows -> 11 segments total
before = db.versions("trades")[-1]
commit = db.compact("trades")
print(f"compacted: {before['segments']} segments -> {commit['segments_total']}, "
      f"rows unchanged: {commit['rows_total']:,}")

[
    {k: v[k] for k in ("sequence", "op", "rows", "segments", "note") if k in v}
    for v in db.versions("trades")
]
output
compacted: 11 segments -> 1, rows unchanged: 141,301
[{'sequence': 0, 'op': 'create', 'rows': 0, 'segments': 0},
 {'sequence': 1,
  'op': 'append',
  'rows': 11418,
  'segments': 1,
  'note': 'day 1: arrow feed'},
 {'sequence': 2,
  'op': 'append',
  'rows': 25773,
  'segments': 2,
  'note': 'day 2: pandas'},
 {'sequence': 3,
  'op': 'append',
  'rows': 38842,
  'segments': 3,
  'note': 'day 3: polars'},
 {'sequence': 4,
  'op': 'append',
  'rows': 49876,
  'segments': 4,
  'note': 'day 4: parquet drop'},
 {'sequence': 5,
  'op': 'append',
  'rows': 63593,
  'segments': 5,
  'note': 'day 5: legacy csv'},
 {'sequence': 6,
  'op': 'append',
  'rows': 76641,
  'segments': 6,
  'note': 'day 6: CAS append'},
 {'sequence': 7,
  'op': 'append',
  'rows': 90447,
  'segments': 7,
  'note': 'daily load 2026-06-09'},
 {'sequence': 8,
  'op': 'append',
  'rows': 102616,
  'segments': 8,
  'note': 'daily load 2026-06-10'},
 {'sequence': 9,
  'op': 'append',
  'rows': 113744,
  'segments': 9,
  'note': 'daily load 2026-06-11'},
 {'sequence': 10,
  'op': 'append',
  'rows': 126837,
  'segments': 10,
  'note': 'daily load 2026-06-12'},
 {'sequence': 11,
  'op': 'append',
  'rows': 141301,
  'segments': 11,
  'note': 'daily load 2026-06-15'},
 {'sequence': 12, 'op': 'compact', 'rows': 141301, 'segments': 1}]
# One tape, five formats, eleven commits - and SQL sees a single clean table.
db.sql(
    """
    SELECT time_bucket('1d', ts) AS session, count(*) AS trades,
           round(sum(price * size) / 1e6, 1) AS notional_mm
    FROM trades GROUP BY session ORDER BY session
    """
).to_pandas()
output
session trades notional_mm
0 2026-06-01 00:00:00+00:00 11418 250.4
1 2026-06-02 00:00:00+00:00 14355 339.6
2 2026-06-03 00:00:00+00:00 13069 295.4
3 2026-06-04 00:00:00+00:00 11034 252.4
4 2026-06-05 00:00:00+00:00 13717 309.5
5 2026-06-08 00:00:00+00:00 13048 294.2
6 2026-06-09 00:00:00+00:00 13806 319.5
7 2026-06-10 00:00:00+00:00 12169 264.3
8 2026-06-11 00:00:00+00:00 11128 251.3
9 2026-06-12 00:00:00+00:00 13093 288.0
10 2026-06-15 00:00:00+00:00 14464 338.7

Takeaways#

db.close()