Skip to content

feat(iceberg): Apache Iceberg export (zero-copy, opt-in) - #532

Merged
xe-nvdk merged 16 commits into
26.09.1from
feat/iceberg-export-eval
Jul 14, 2026
Merged

xe-nvdk merged 16 commits into
26.09.1from
feat/iceberg-export-eval

Conversation

@xe-nvdk

@xe-nvdk xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member

Summary

Adds an opt-in Apache Iceberg export layer: Arc publishes its existing Parquet files as standard Iceberg tables so any Iceberg-aware engine (Spark, Trino, DuckDB, Snowflake, PyIceberg) can query Arc's data directly — no proprietary API in the path. This is the strongest form of the "your data, no lock-in" promise.

Key property: zero-copy. A background reconciler registers Arc's existing Parquet files into an Iceberg table by reference (Iceberg add_files — no data rewrite) and keeps the file set in sync as compaction/retention change the underlying files. The 20M+ rec/s ingest path is completely unchanged.

iceberg-go v0.6.0 is a normal published dependency — no replace directives, CI-safe.

What's included

  • Reconciler core — diffs the on-disk Parquet set against the Iceberg table's live files and commits the delta via ReplaceDataFiles (register/deregister by path, no rewrite). Idempotent + self-healing (driven by durable storage, not an event stream).
  • Config + wiring — iceberg.enabled (default false); wired outside the cluster block so it runs in OSS; writer-gated in cluster mode via the compaction gate.
  • Directory-reader discovery — emits version-hint.text + v<N>.metadata.json so DuckDB/Spark can read from the table directory (both verified).
  • Snapshot expiry + metadata pruning — bounds snapshot history and metadata growth (iceberg.retain_snapshots, default 10).
  • Schema evolution — measurements that gain columns evolve the Iceberg table automatically; older files stay readable.
  • Backup/restore — Iceberg warehouse metadata is included in Arc's backup (previously restore would lose the tables).
  • Incremental reconcile — skips measurements whose file set is unchanged since the last pass (no footer reads/diff at steady state) + a per-measurement timeout so a slow measurement can't block shutdown.

Also fixes a pre-existing latent bug: cfg.Backup was never populated in config.Load(), so the backup/restore API never enabled regardless of the [backup] TOML section (verified: startup now logs "Backup/restore enabled", GET /api/v1/backup/ returns 200).

Cross-engine verification

The Iceberg tables Arc produces (including the schema-evolution case) were read correctly by four independent engines: DuckDB (C++), PyIceberg (Python reference impl), Apache Spark on host (JVM), and Apache Spark in a container (JVM, isolated). Name-mapping resolves Arc's field-ID-less Parquet on all of them.

Limitations (v1) — enforced at config load

  • Local storage only. iceberg.enabled=true requires storage.backend="local" and is incompatible with cold-tier tiering (a file migrated to S3 would leave the Iceberg table). Arc refuses to start with either combination.
  • Eventual consistency — the Iceberg view lags live ingest by up to reconcile_interval.
  • Cluster — exactly one node must run the reconciler; guaranteed under the compactor failover lease, documented for static-role clusters.

Review

This went through the CLAUDE.md pre-PR discipline:

  • Config matrix (7 rows: OSS / OSS+compaction / cluster / cluster+failover / cluster+static-roles / +cold-tier / +non-local-backend) — traced from startup through the actual if-blocks. It surfaced a real gap (guard missed a non-local primary backend), now fixed + tested.
  • Single deep reviewer with the matrix — confirmed every row, ran the five checks (precondition trace, OSS smoke, failure modes, doc-vs-code drift, hot-path/SQLite checklist). Zero nil-derefs, zero doc drift, all guards in place.

Test plan

  • go build ./cmd/... ./internal/...
  • go test ./internal/iceberg/ ./internal/config/ ./internal/backup/ (unit + integration, incl. schema evolution, expiry+pruning, incremental skip, day-straddle skip, backup/restore round-trip, all config guards)
  • go test -race ./internal/iceberg/...
  • gofmt -l clean, go vet clean
  • Binary run: ingest → reconcile → external read (DuckDB/PyIceberg/Spark), incl. schema-evolved column and post-expiry read
  • Gemini review

Docs

User guide added in docs.basekick.net (Integrations → Apache Iceberg); release notes updated (RELEASE_NOTES_2026.09.1.md).

🤖 Generated with Claude Code

Ignacio Van Droogenbroeck and others added 9 commits July 13, 2026 17:18
New internal/iceberg package that projects Arc's existing Parquet files into
Apache Iceberg tables, so external engines (Spark/Trino/DuckDB) can read Arc's
data without changing the ingest path. Pure library; not yet wired into main.go.

- Exporter over a SQLite Iceberg catalog, reusing Arc's existing mattn/go-sqlite3
  driver (Arc is already a CGo build via duckdb-go + mattn — no new SQLite impl).
- EnsureTable: day(time) partition spec (daily-compacted files span 24h, so hour()
  would reject them) + time->timestamptz mapping (Arc writes UTC-adjusted
  TIMESTAMP_MICROS; plain `timestamp` fails AddFiles).
- ReconcileMeasurement: diffs Arc's durable file set vs the table's live files
  (Scan().PlanFiles) and commits the delta via ReplaceDataFiles in one snapshot.
  Idempotent (no-op when converged) and driven by durable state, so a missed/failed
  commit self-heals next tick — unlike an event-stream design.
- SchemaFromParquet: pure-Go footer read -> Iceberg schema (no DuckDB handle).
- PathResolver: storage-relative key -> file://|s3://|azure:// URI.

Verified: iceberg-go v0.6.0 AddFiles auto-emits schema.name-mapping.default so
Arc's field-ID-less Parquet reads correctly from a non-DuckDB engine (PyIceberg,
Phase 0b). Test covers add + idempotency + hourly->daily supersession.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…se 2)

Wires the Iceberg export reconciler into Arc as an opt-in feature
(storage.file_format unaffected; iceberg.enabled=false by default). Verified
end-to-end in the running binary: ingest -> flush -> reconcile -> external
DuckDB read, including a column added mid-stream.

- IcebergConfig (config + setDefaults + ARC_ICEBERG_* env); opt-in.
- FileSetSource = StorageWalkSource: lists measurements/files straight from the
  storage backend, so it works in OSS and cluster with no tiering/Raft dependency
  (tier_files is only populated when tiering is enabled — would be empty in OSS).
- Scheduler: interval loop, writer-gated (reuses the compaction cluster gate;
  nil/always-run in OSS), converges each pass (idempotent reconciler).
- main.go: constructed OUTSIDE the cluster block so it runs standalone; clean
  shutdown hook; warehouse defaults to the storage root.

Two bugs the running binary caught that unit tests missed, both fixed + regression
tested:
- Schema evolution: Arc's schema-flexible ingest yields files with differing
  columns; a wider file failed AddFiles. Fixed via UnionSchema (superset across a
  measurement's files) + evolveSchema (AddColumn as optional).
- Name-mapping not extended on evolution: UpdateSchema.AddColumn leaves the
  evolution-added column out of schema.name-mapping.default, so external readers
  (Arc's Parquet has no field IDs) fail with "no field-mapping exists". Fixed by
  refreshing the mapping from the actual post-evolution schema.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…ased readers

Publishes Hadoop-catalog discovery files next to each table's metadata so readers
that point at the table DIRECTORY (no exact metadata filename, no catalog) can find
the current metadata:

- metadata/version-hint.text — the version integer, NO trailing newline (DuckDB
  reads it verbatim; a newline makes it seek "v<N>\n.metadata.json").
- metadata/v<N>.metadata.json — a copy of the current metadata under the
  Hadoop-convention name. Required IN ADDITION to the hint: iceberg-go/the SQL
  catalog write "NNNNN-<uuid>.metadata.json", but Spark and DuckDB resolve the hint
  strictly to "v<N>.metadata.json" (both fail "metadata file for version N missing"
  with only the hint).

Written after every commit (create/reconcile/evolve) via the storage backend, so it
works for local and S3 warehouses. Best-effort: failures are logged, not fatal — the
SQL catalog stays the source of truth and the next reconcile rewrites these.
Catalog-aware readers (PyIceberg, iceberg-go) are unaffected.

Verified against Arc's binary + both directory-based engines (previously failing):
DuckDB iceberg_scan('<tableDir>') and Spark read.format("iceberg").load('<tableDir>')
now read correctly. NewExporter gains a storage.Backend param.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Bounds snapshot history and metadata growth so Iceberg tables don't accumulate
unbounded state. New config iceberg.retain_snapshots (default 10).

- After each reconcile commit, ExpireSnapshots(WithRetainLast(retain),
  WithOlderThan(0)). WithRetainLast alone is a FLOOR, not a cap: iceberg-go only
  expires a snapshot older than maxSnapshotAgeMs (default ~5d) AND beyond the
  retain count, so with the default age nothing expires. WithOlderThan(0)
  satisfies the age gate so retain-last becomes the effective cap.
- Table created with write.metadata.delete-after-commit.enabled=true +
  write.metadata.previous-versions-max=retain, so iceberg-go prunes its own
  NNNNN-*.metadata.json history.
- pruneOldVersionFiles keeps only the newest `retain` of OUR v<N>.metadata.json
  copies (iceberg-go doesn't know about those). Scan-based (lists the dir), robust
  to non-contiguous versions — each reconcile commits twice (ReplaceDataFiles +
  ExpireSnapshots) so versions skip. Never deletes the current version.

All best-effort: an expiry/prune failure doesn't undo the reconcile; the next pass
retries. Writer-gated (single writer), so no reader-vs-expire race beyond Iceberg's
snapshot isolation.

Verified in the binary: 8 distinct writes -> live snapshot count capped at 3
(retain=3), v<N> copies capped at 3, and DuckDB dir-scan still reads all 8 rows
after expiry. NewExporter gains a retain param.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…hase 5)

CreateBackup previously copied only .parquet files, silently dropping the Iceberg
warehouse metadata (metadata.json / .avro manifests / version-hint.text) — a restore
kept the parquet data but lost the Iceberg tables pointing at it. Now those metadata
files are copied via the same mechanism (isIcebergMetadata predicate: files under a
"/metadata/" dir ending .metadata.json or .avro, or version-hint.text), kept out of
the db/measurement inventory. Restore needs no change: restoreDataFiles round-trips
everything under {backupID}/data/ by path, so the metadata returns to its original
location and the SQLite catalog's pointers resolve.

Verified (TestBackupRestore_IcebergMetadata): backup -> delete metadata from the data
store -> restore -> all Iceberg metadata files return with correct content + the
parquet. Predicate has a unit test incl. the false-positive guard (a measurement named
"metadata" holding parquet must not match).

Note: found a pre-existing, unrelated gap — cfg.Backup is never populated in
config.go Load(), so the backup API is not enabled regardless of the [backup] TOML
section. The Iceberg backup change is correct + tested at the Manager level; that
wiring gap deserves its own fix. (Documented in docs/progress.)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Adversarial review of the Iceberg export feature found two blockers and several
high/medium issues, all verified against iceberg-go v0.6.0 source and fixed:

B1 (SQLite): the catalog opened the shared SQLite file with no pragmas, unlike auth
(WAL + busy_timeout + SetMaxOpenConns(1)). Under concurrent load on the shared DB
(auth/audit/tiering/retention/MQTT) the catalog got immediate SQLITE_BUSY, failing
reconciles. Now opened with ?_journal_mode=WAL&_busy_timeout=5000&_foreign_keys=ON
and SetMaxOpenConns(1), mirroring auth. (cmd/arc/main.go)

B2 (partition panic): iceberg-go panics->errors "more than one value for partition
field" when a single Parquet file's time min/max straddle a UTC-day boundary (rare:
backfill / midnight flush), which wedged the whole measurement's export forever.
replaceDataFilesResilient now tries the batch, and on that specific error falls back
to per-file adds, SKIPPING only the straddling file(s) (logged at Error) so the rest
still export. Test with a 2-day-spanning file. (exporter.go)

H1 (self-referential walk): the warehouse (arc_<db>.db/) lives under the storage
root, so Measurements() enumerated it as a phantom database. NewStorageWalkSource now
takes nsPrefix and skips <nsPrefix>_*.db dirs. New source_test.go covers this + that
Files() recurses nested Y/M/D/H partitions.

H2 (cold tiering): a file migrated to a cold S3 tier vanishes from the local walk and
would be DELETED from the Iceberg table. config.Load now rejects
iceberg.enabled + tiered_storage.cold.enabled. Test added.

Also: H3 — documented the single-compaction-node requirement for Phase-4 clusters
(non-transactional warehouse writes). M1 — ExpireSnapshots failures now log at Error
(persistent failure grows history unbounded). L3 — fixed package-doc drift
(ReplaceDataFiles, not AddFiles+Delete).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…se notes

Deferred optimizations from the adversarial review:

- Incremental reconcile: the scheduler now fingerprints each measurement's file set
  (data files + local files) and SKIPS measurements whose set is unchanged since the
  last successful pass — no footer reads, no diff, no commit. UnionSchema (the O(files)
  footer scan) runs only when the set actually changed. At steady state a pass is cheap
  regardless of total file count. Stale cache entries for deleted measurements are pruned.
- Per-measurement timeout (2m) so one slow/wedged measurement can't block the whole pass
  or graceful shutdown; ctx cancellation is checked between measurements.

Pass log now reports reconciled/unchanged/failed. Test: unchanged pass creates no new
snapshot; adding a file triggers a reconcile.

Also adds RELEASE_NOTES_2026.09.1.md coverage for the Iceberg export feature
(what/why/how, config table, limits). User guide added in docs.basekick.net
(Integrations -> Apache Iceberg).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
… enables

BackupConfig had a struct + setDefaults entries (backup.enabled=true,
backup.local_path) but was never read in config.Load(), so cfg.Backup was always
the zero value — cfg.Backup.Enabled=false — and the backup/restore API
(POST /api/v1/backup) was never registered regardless of the [backup] TOML section.
Pre-existing bug, found while testing the Iceberg backup work.

Load() now populates Backup from the "backup.enabled" / "backup.local_path" keys,
mirroring the other optional-feature blocks. Verified against the binary: startup now
logs "Backup/restore enabled" and GET /api/v1/backup/ returns 200 (was 404). Test
asserts the defaults land and env overrides apply.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
The pre-PR config matrix surfaced a gap: the guard refused iceberg + cold-tier
tiering but NOT iceberg + a non-local PRIMARY backend (S3/Azure). Iceberg export is
local-only in v1 (the reconciler walks the single backend and reads Parquet footers
from the local filesystem; verified only against local), and the docs/release notes
say so — but nothing enforced it, so a primary-S3 deployment would run the reconciler
against S3 unguarded.

config.Load now also requires storage.backend="local" when iceberg.enabled=true,
alongside the existing cold-tier refusal. Both fire before startup. Test added
(TestLoad_IcebergRequiresLocalBackend); the deep review confirmed the guard is correct
and breaks no existing valid config (only fires when iceberg.enabled, default false).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer that periodically reconciles Arc's Parquet files into Iceberg tables, along with configuration updates, backup support for Iceberg metadata, and environment validation. The review feedback highlights several critical improvements: normalizing Windows file paths to prevent malformed URIs, fixing snapshot expiration by using the current timestamp instead of zero, and optimizing schema derivation incrementally to avoid costly sequential reads of Parquet footers.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/paths.go
Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/scheduler.go
Comment thread internal/iceberg/scheduler.go
Comment thread internal/iceberg/schema.go

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer that periodically reconciles Arc's Parquet files with Iceberg tables, and updates the backup and restore system to preserve Iceberg metadata. Key feedback from the review highlights several critical improvements: ensuring robust name-mapping self-healing during schema evolution, fixing a bug where snapshot expiration was a no-op due to an incorrect age gate, and converting relative file:// URIs to absolute paths with forward slashes for cross-platform compatibility. Additionally, the review suggests optimizing the scheduler by replacing the ticker with a timer to prevent back-to-back runs, avoiding redundant storage listings, and caching Parquet schemas to prevent expensive O(N) footer reads.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/paths.go Outdated
Comment thread internal/iceberg/paths.go
Comment thread internal/iceberg/scheduler.go
Comment thread internal/iceberg/scheduler.go Outdated
Comment thread internal/iceberg/schema.go Outdated
…ng self-heal, incremental schema

- file:// URIs now absolute + forward-slashed (Windows-safe) via localFileURI
- heal schema.name-mapping.default if a prior evolution's mapping txn failed
- incremental schema derivation: read only newly-added footers, MergeSchemas
- scheduler: Ticker→Timer so a slow pass can't cause back-to-back runs
- single storage List per measurement (FilesAndLocal)
- tests: incremental schema evolution + MergeSchemas conflict

Declined the snapshot-expiry finding: WithOlderThan takes a time.Duration,
not Unix millis; WithOlderThan(0) is correct (verified: 5 snaps, retain=3 → 3
kept). Gemini's suggested time.Now().UnixMilli() would disable expiry.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Thanks for the review — pushed 60436d8 addressing 5 of the 6 findings.

Fixed:

  • file:// URIs (paths.go) — now absolute + forward-slashed via a localFileURI helper (filepath.Abs + filepath.ToSlash, plus the Windows file:///C:/… drive case). Verified in a running instance: table location is now file:///….
  • Name-mapping self-heal (exporter.go) — added healNameMapping; when a prior evolution committed the schema but its follow-up mapping txn failed, len(missing)==0 no longer skips repair. Commits only when the stored mapping actually diverges.
  • Incremental schema derivation (scheduler.go + MergeSchemas in schema.go) — a grown measurement now reads footers of only the newly-added local files and merges into the cached schema, instead of an O(N) re-read of every footer. Cached localFiles set in measurementState. New tests: TestScheduler_IncrementalSchemaEvolution, TestMergeSchemas.
  • time.Ticker → time.Timer (scheduler.go) — reset after each pass guarantees a full idle interval; a slow pass can no longer queue back-to-back runs.
  • Redundant List (source.go) — FilesAndLocal lists the prefix once and returns both the URI refs and local paths; Files/LocalFiles delegate.

Declined — snapshot expiry (exporter.go:355, WithOlderThan(0)): this one is a false positive, and the suggested fix would disable expiry. In iceberg-go v0.6.0, WithOlderThan takes a time.Duration (a max-age), not a Unix-millis timestamp:

func WithOlderThan(t time.Duration) ExpireSnapshotsOpt {
    return func(cfg *expireSnapshotsCfg) { n := t.Milliseconds(); cfg.maxSnapshotAgeMs = &n }
}

The expiry loop keeps a snapshot while !(snapAge > maxSnapshotAgeMs && numSnapshots >= minSnapshotsToKeep). With WithOlderThan(0) → maxSnapshotAgeMs = 0, the age gate snapAge > 0 is always true, so WithRetainLast(retain) becomes the effective cap ("keep the last N, expire the rest") — exactly the intent.

The suggestion WithOlderThan(time.Now().UnixMilli()) passes an int64 interpreted as time.Duration nanoseconds → maxSnapshotAgeMs ≈ 1.78×10⁶ (~29 min), so a fresh snapshot has snapAge (~30s) > 29min = false → the age gate never passes → nothing ever expires, reintroducing the unbounded-growth bug.

Verified empirically against the running binary: with retain_snapshots = 3, ingesting 5 distinct files produced 5 append snapshots and the metadata retained exactly 3 — expiry works as written.

@gemini-code-assist please review the follow-up commit.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer for Arc, enabling external query engines to directly query Arc's Parquet data. It implements a background reconciler (Exporter and Scheduler) that periodically walks the storage backend, derives schemas from Parquet footers, and registers/deregisters files in Iceberg tables. Additionally, the backup and restore processes are updated to include Iceberg metadata, configuration validation is added to ensure Iceberg is local-only and incompatible with cold-tier storage, and comprehensive tests are introduced. The reviewer feedback suggests parallelizing Parquet footer reads in UnionSchema to prevent timeouts on large datasets, and recommends using path/filepath instead of path for better Windows compatibility.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/schema.go Outdated
Comment thread internal/iceberg/source.go Outdated
Comment thread internal/iceberg/source.go Outdated
Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/exporter.go Outdated
Comment thread internal/iceberg/schema.go
…th.Base for Windows keys

- UnionSchema now reads Parquet footers in parallel (errgroup, limit 8) so first-sight
  full derivation of a many-file measurement can't exceed measurementTimeout. Merge stays
  deterministic (results collected by index, folded in path order). x/sync already a dep.
- path.Base → filepath.Base at the two key-parsing sites (isDataFile, pruneOldVersionFiles/
  parseVersionAndMetaDir): LocalBackend.List returns backslash keys on Windows (filepath.Rel,
  no ToSlash) and Iceberg export is local-only. path.Join/path.Dir kept as-is — they BUILD
  forward-slash storage keys written back to the backend, where filepath would corrupt them.
- test: TestUnionSchema_ParallelDeterministic (order-stable under concurrency + error surfaced).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 2 addressed in 35b045e — both findings taken.

Parallel footer reads (UnionSchema) — now reads footers concurrently via errgroup (limit 8; golang.org/x/sync is already a dependency). This only matters on first-sight/full derivation of a many-file measurement, since round 1's incremental path already restricts steady-state passes to reading only newly-added footers — but the full-derivation case is exactly where the measurementTimeout risk lived, so it's worth it. The merge stays deterministic: results are collected by index and folded in localPaths order, so column ordering and conflict detection are unchanged. Added TestUnionSchema_ParallelDeterministic (order-stable under concurrency + per-file error surfaced). Passes under -race.

path.Base → filepath.Base — applied at the two sites that parse a filename out of an incoming key (isDataFile, parseVersionAndMetaDir). Good catch: LocalBackend.List returns keys via filepath.Rel with no ToSlash, so on Windows they're backslash-separated, and Iceberg export is local-only — so these keys really can carry \.

One deliberate exception: I kept path.Join/path.Dir at the other sites in exporter.go (building dirKey, v<N>.metadata.json keys, version-hint.text key). Those construct storage keys that are written back to the backend, where the object-key contract is forward-slash on every backend — switching them to filepath would emit backslash keys on Windows and corrupt the warehouse layout. So the rule applied is: filepath.Base when parsing an incoming (possibly-OS) key's filename; path.* when building a forward-slash storage key.

@gemini-code-assist please review the follow-up.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export feature that publishes Arc's existing Parquet files as Iceberg tables without rewriting data. It includes a background reconciler, schema evolution support, and backup/restore integration for Iceberg metadata. A critical issue was identified in the scheduler where deleting all files of a measurement would skip the reconciliation pass entirely, leaving orphaned references in the Iceberg catalog; a fix was suggested to allow reconciling empty tables.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/scheduler.go
…d (Gemini round 3)

Retention/compaction deleting every file of a measurement left the Iceberg table
pointing at deleted paths forever: Arc's Delete only os.Remove()s the file, so the
{db}/{measurement}/... dir tree survives, Measurements() keeps yielding it, and
reconcileOne's len(localFiles)==0 early-return skipped the pass. External engines
then fail on the missing files. Reproduced, then fixed: files=[] now reconciles the
table to empty.

Two guards beyond the reported finding:
- Only empty a table that ALREADY exists (TableExists). A stray empty directory that
  never held data must not mint a zero-column table — EnsureTable with an empty schema
  would create one with no fields and no partition spec.
- Fingerprint-cache the empty state so a permanently-empty measurement is emptied ONCE
  and skipped thereafter, instead of re-committing and re-logging every tick. The cached
  state has a nil localFiles, so a later re-ingest falls through to full re-derivation.

Verified on the running binary: 40 rows -> delete all -> exactly one "reconciled to
empty" across 5 ticks -> DuckDB iceberg_scan returns 0 rows (was: read failure).

Tests: TestScheduler_AllFilesDeletedEmptiesTable (empty-out + no-repeat + re-ingest
recovery), TestScheduler_EmptyDirNeverCreatesTable.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 3 addressed in 01de37c — good catch, this was a real data-correctness bug and I reproduced it before fixing.

Confirmed the mechanism: LocalBackend.Delete only os.Removes the file with no empty-directory pruning, and ListDirectories uses os.ReadDir — so after retention deletes every file, the {db}/{measurement}/… tree survives, Measurements() keeps yielding the measurement, and reconcileOne saw files=[]/localFiles=[] and skipped. Repro before the fix:

after first pass: table has 1 files
Measurements() after deleting all files: [{mydb cpu}]
after delete+reconcile: table STILL has 1 files (want 0)   <-- orphaned ref

Fixed, plus two guards the suggested patch would have missed:

  1. Only empty a table that already exists (TableExists). The literal len(files) > 0 && len(localFiles) == 0 change lets a stray empty directory that never held data fall through to ReconcileMeasurement → EnsureTable with an empty schema, which mints a table with zero columns and no partition spec. Guarded + covered by TestScheduler_EmptyDirNeverCreatesTable.

  2. Fingerprint-cache the empty state. My first cut dropped the cache entry on empty-out, which made every subsequent tick re-run the catalog load + commit and re-log — I caught it on the binary (the message fired twice in 8s). Now the empty state is fingerprinted like any other, so it empties once and reports unchanged after. The cached empty state has a nil localFiles, so a later re-ingest correctly falls through to full re-derivation rather than the incremental branch.

Verified on the running binary: 40 rows → delete all files → exactly one reconciled to empty across 5 reconcile ticks → iceberg_scan returns 0 rows (previously: read failure on missing files). Tests added for empty-out, no-repeat, and re-ingest recovery. All green under -race.

@gemini-code-assist please review the follow-up.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer to Arc, allowing its Parquet data files to be registered as Iceberg tables and queried directly by external engines. It includes a background reconciler, configuration options, and updates to the backup manager to ensure Iceberg metadata is preserved during backup and restore operations. The review feedback suggests several robust improvements: using the idiomatic Equals method for Iceberg type comparisons, extending type mapping to support 32-bit integers and floats, ensuring S3 prefixes are resolved with a trailing slash to prevent malformed URIs, and future-proofing Iceberg metadata detection in backups by matching any non-parquet file under the metadata directory.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/schema.go Outdated
Comment thread internal/iceberg/schema.go
Comment thread internal/iceberg/schema.go
Comment thread internal/iceberg/paths.go
Comment thread internal/backup/backup.go
… metadata backup (Gemini round 4)

- schema: compare Iceberg types with Type.Equals() instead of String() comparison
  (idiomatic; the interface exposes Equals(Type) bool).
- schema: map arrow Int32/Float32 -> iceberg int/float. Arc's own ingest only emits
  Int64/Float64/String/Boolean/Timestamp_us, but Parquet also arrives via the bulk
  import path (internal/api/import.go) carrying externally-produced 32-bit columns —
  those previously failed the measurement's whole export with "unsupported Arrow type".
  Same reason Decimal128 was already mapped.
- backup: isIcebergMetadata is now a catch-all (any non-parquet under a /metadata/
  segment) instead of an extension allowlist, so new Iceberg metadata types (e.g. Puffin
  .puffin stats/index files) can't silently drop out of the backup and be lost on restore.
  Safe because CreateBackup's switch tests the .parquet branch first, so data files never
  reach the predicate.
- tests: TestArrowToIceberg (full mapping incl. timestamptz vs timestamp + unsupported
  type error); isIcebergMetadata cases for .puffin and stray non-parquet.

Declined the S3 trailing-slash finding: S3Backend.prefix is sanitized at construction via
SanitizeS3Prefix ("" or guaranteed trailing /), so the malformed-URI case can't occur — and
the S3 arm is unreachable anyway (Iceberg export refuses non-local backends).

Verified on the running binary: backup captured all 7 Iceberg metadata files (metadata.json,
.avro manifest+snapshot, version-hint.text) alongside the parquet data.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 4 addressed in 404233a — 3 of 4 taken.

Fixed:

  • Type.Equals() instead of String() comparison (UnionSchema + MergeSchemas) — confirmed Equals(Type) bool is on the iceberg.Type interface. Agreed, the string compare was fragile.
  • int32/float32 mapping. Worth noting why this is real: Arc's own ingest writer only ever emits Int64/Float64/String/Boolean/Timestamp_us, so at first glance these look unreachable — but Parquet also enters storage via the bulk import path (internal/api/import.go), which can carry externally-produced 32-bit columns. Those previously failed the measurement's entire export with "unsupported Arrow type". Same reason Decimal128 was already mapped. Added TestArrowToIceberg covering the full table (including timestamptz-vs-timestamp and the unsupported-type error).
  • Backup metadata catch-all. Good future-proofing call — Puffin files are exactly the case that would silently drop out and be lost on restore. Safe as suggested because CreateBackup's switch tests the .parquet branch before this predicate, so data files never reach it (a measurement literally named metadata still can't be misclassified — covered by test). Verified on the running binary: backup captured all 7 Iceberg metadata files (both .metadata.json forms, .avro manifest + snapshot, version-hint.text) alongside the parquet, with the manifest inventory still counting only the data file.

Declined — S3 prefix trailing slash (paths.go:63): the invariant is already guaranteed upstream. S3Backend.prefix is sanitized at construction:

// s3.go:177-178
// Sanitize prefix: strip leading /, reject .., ensure trailing / if non-empty
prefix := SanitizeS3Prefix(cfg.Prefix)

and the field is declared prefix string // path prefix within the bucket (sanitized, with trailing /). SanitizeS3Prefix returns either "" or a prefix ending in /, so s3://bucket/prefixmydb/... can't occur. Adding a second guard would imply the constructor's contract isn't trustworthy. It's also moot in v1 regardless: Iceberg export refuses to start on a non-local backend, so the S3 arm is unreachable — if/when object-store warehouses land, that's the point to re-validate this path end-to-end rather than defensively now.

@gemini-code-assist please review the follow-up.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer to Arc, allowing existing Parquet files to be published as Iceberg tables via a periodic reconciler. It includes changes to configuration, backup/restore logic to preserve Iceberg metadata, and the core implementation of the exporter and scheduler under internal/iceberg/. Feedback from the reviewer suggests several robustness improvements: passing a context.Context to UnionSchema to support early cancellation during parallelized file I/O (along with updating its callers and tests), adding type compatibility checks during schema evolution to prevent table corruption, checking for "null" name mappings in healNameMapping, and refining isAlreadyExists to avoid false positives on errors containing "not exist".

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/schema.go Outdated
Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/scheduler.go Outdated
Comment thread internal/iceberg/exporter_test.go Outdated
Comment thread internal/iceberg/scheduler_test.go Outdated
Comment thread internal/iceberg/exporter.go Outdated
Comment thread internal/iceberg/exporter.go
…d already-exists (Gemini round 5)

- UnionSchema now takes a context and uses errgroup.WithContext: the parallel footer
  reads abort on the first failure or on reconciler timeout/shutdown instead of
  grinding through every remaining file. Callers/tests updated.

- evolveSchema now rejects a column whose type CHANGED on an existing table. This was a
  real silent-corruption gap: UnionSchema only compares the current pass's files against
  each other, and MergeSchemas compares against the in-memory cache — which is empty after
  a restart or an empty-out. So `value` long -> all old files age out -> restart -> new
  files arrive with value as double would derive a self-consistent schema, see the column
  as "present", skip it, and register type-incompatible files into the table. Now the
  measurement fails loudly and the table stays readable.

- healNameMapping: never write schema.name-mapping.default="null". NameMapping is a
  []MappedField, so a nil mapping marshals to the literal "null" — an invalid mapping that
  would break the external readers the property exists to serve.

- isAlreadyExists: match iceberg-go's typed sentinels via errors.Is (the SQL catalog wraps
  them with %w) instead of substring-matching "exist", which also matched "does not exist"
  and could swallow a genuine not-found/backend error as success.

Tests: UnionSchema cancellation; evolveSchema type-mismatch rejected + unchanged-schema
no-op. Verified on the running binary: two measurements in one namespace both create and
reconcile (exercises the already-exists path), DuckDB reads both.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 5 addressed in f920460 — all four taken, one implemented differently.

  • ctx in UnionSchema — now UnionSchema(ctx, paths) with errgroup.WithContext, so the parallel footer reads abort on the first failure or on the reconciler's measurementTimeout/shutdown instead of reading every remaining file. All callers and tests updated. Added TestUnionSchema_ContextCancelled.

  • Schema-evolution type check — this was the important one. You're right, and the gap is worse than it looks from the diff: neither derivation path can catch it. UnionSchema only compares the current pass's files against each other, and MergeSchemas compares against the in-memory cache, which is empty after a restart or an empty-out. So the concrete failure is: value is long → retention ages out every old file → Arc restarts (cache gone) → new files arrive with value as double → full derivation succeeds (all files agree), evolveSchema sees the column as "present" and skips it → type-incompatible files get registered. evolveSchema is the only place that sees the durable table schema, so that's the right guard. Now it fails the measurement loudly and the table stays readable. Covered by TestEvolveSchema_TypeMismatchRejected (plus an unchanged-schema no-op assertion so the check can't false-positive).

  • "null" name-mapping — confirmed mechanically: NameMapping is []MappedField, so a nil mapping marshals to the literal "null", and writing schema.name-mapping.default="null" would break the very readers the property exists for. Guarded (warn + skip, leaving it unset rather than invalid).

  • isAlreadyExists — agreed the old Contains(msg, "exist") also matched "does not exist". I went further than the suggested string patch: iceberg-go exposes typed sentinels (catalog.ErrNamespaceAlreadyExists / ErrTableAlreadyExists) and the SQL catalog wraps them with %w, so errors.Is is exact and immune to message-wording changes:

func isAlreadyExists(err error) bool {
	return errors.Is(err, icecatalog.ErrNamespaceAlreadyExists) ||
		errors.Is(err, icecatalog.ErrTableAlreadyExists)
}

Verified on the running binary (this path runs on every EnsureTable, so a regression would break table creation outright): ingested two measurements into one database — both cpu and mem tables created and reconciled in the same namespace, exercising the already-exists branch on the second, and iceberg_scan returns 30 rows from each.

@gemini-code-assist please review the follow-up.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer that publishes Arc's existing Parquet files as Iceberg tables via a background reconciler. It includes schema derivation and evolution, a periodic scheduler, backup/restore support for Iceberg metadata, and configuration constraints. Feedback on the implementation highlights two issues in internal/iceberg/exporter.go: first, using WithOlderThan(0) in expireSnapshots prevents snapshots from ever expiring because they are not older than the Unix epoch; second, marshaling a nil name mapping in evolvedSchema can write an invalid "null" property value to the table, which would break external readers.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/exporter.go
Comment thread internal/iceberg/exporter.go Outdated
…ni round 6)

Round 5 guarded healNameMapping against writing schema.name-mapping.default="null"
(NameMapping is a []MappedField, so a nil mapping marshals to the literal "null" — an
invalid mapping that breaks the external readers the property exists to serve), but the
sibling site in evolveSchema did the same marshal+SetProperties WITHOUT the guard.

Rather than duplicate the check, both paths now route through one helper,
setNameMappingFromSchema — the only place that property is set. The two sites doing the
same thing and drifting apart is exactly how this was missed, so the guard now lives in
a single place by construction. healNameMapping keeps its best-effort semantics (warn,
retry next pass) by wrapping the shared helper.

Verified on the running binary: narrow -> wide ingest evolves the schema
("added_columns":1), the resulting name-mapping includes the evolution-added column, and
DuckDB reads the evolved table (60 rows, humidity populated in 30, NULL in the older 30).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 6 — one taken, one is a repeat of a finding already disproved in round 1.

Taken: null name-mapping in evolveSchema (4393b9c). Correct, and a genuine miss on my part — round 5 guarded healNameMapping against writing schema.name-mapping.default="null", but the sibling site in evolveSchema did the same marshal + SetProperties without the guard. Rather than duplicate the check, both paths now route through one helper (setNameMappingFromSchema) that is the only place the property is set — the two sites doing the same thing and drifting apart is precisely how this was missed, so the guard now lives in a single place by construction.

Verified on the running binary: narrow → wide ingest evolves the schema ("added_columns":1), the resulting name-mapping includes the evolution-added column, and DuckDB reads the evolved table (60 rows, humidity populated in 30, NULL in the older 30).


Declined (repeat): WithOlderThan(0) in expireSnapshots. This is the same finding from the round-1 review; the detailed refutation is here. Briefly, because it keeps recurring:

iceberg-go v0.6.0's WithOlderThan takes a time.Duration (a max-age), not a Unix-millis timestamp:

func WithOlderThan(t time.Duration) ExpireSnapshotsOpt {
    return func(cfg *expireSnapshotsCfg) { n := t.Milliseconds(); cfg.maxSnapshotAgeMs = &n }
}

The keep-loop retains a snapshot while !(snapAge > maxSnapshotAgeMs && numSnapshots >= minSnapshotsToKeep). WithOlderThan(0) → maxSnapshotAgeMs = 0 → the age gate snapAge > 0 is always true, so WithRetainLast(n) becomes the effective cap. That is the intent: keep the last N, expire the rest.

The suggested WithOlderThan(time.Now().UnixMilli()) passes an int64 that Go interprets as a time.Duration in nanoseconds ≈ 29 minutes, so maxSnapshotAgeMs ≈ 1.78e6 and a fresh snapshot (snapAge ≈ 30s) fails the gate → nothing would ever expire. It would introduce the exact bug the comment describes.

Empirically confirmed twice on the running binary: with retain_snapshots = 3, ingesting 5 distinct files produces 5 append snapshots and the metadata retains exactly 3. Expiry works as written.

@gemini-code-assist please review the follow-up.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces an opt-in Apache Iceberg export layer to Arc, allowing external engines to query Arc's Parquet data directly. It implements a background reconciler that periodically walks the storage backend, derives schemas from Parquet footers, handles schema evolution, and registers files into Iceberg tables. Additionally, the backup and restore manager has been updated to include Iceberg metadata. The review feedback is highly constructive, pointing out a redundant prefix check in the data file identification logic that can be simplified.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread internal/iceberg/source.go Outdated
// export is local-only. filepath.Base handles both separators; path.Base would treat a
// backslash key as one component.
base := filepath.Base(p)
if strings.HasPrefix(base, ".") || strings.HasPrefix(base, ".tmp.") {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The condition strings.HasPrefix(base, ".tmp.") is redundant because any string starting with ".tmp." also starts with ".", which is already checked by the first condition strings.HasPrefix(base, "."). You can simplify this check to only check for the "." prefix.

Suggested change
if strings.HasPrefix(base, ".") || strings.HasPrefix(base, ".tmp.") {
if strings.HasPrefix(base, ".") {

…emini round 7)

strings.HasPrefix(base, ".tmp.") was unreachable — anything starting with ".tmp."
already starts with ".", which the preceding check catches. Same behaviour, one
condition. Kept the ".tmp." rationale in the comment, since skipping Arc's in-flight
writes is the reason the dotfile check exists at all.

Added TestIsDataFile to lock the behaviour in (it was untested): .tmp.* and other
dotfiles skipped, plain/compacted .parquet accepted, non-parquet rejected.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@xe-nvdk

xe-nvdk commented Jul 14, 2026

Copy link
Copy Markdown
Member Author

Round 7 addressed in a25494c — taken, you're right that it's dead code. Anything starting with ".tmp." already starts with ".", so the second HasPrefix was unreachable. Simplified to the single check, keeping the .tmp. rationale in a comment since skipping Arc's in-flight writes is why the dotfile check exists (registering a partially-written Parquet would hand external readers a corrupt file).

Also added TestIsDataFile while I was in there — the predicate was untested, and I did not want to touch it without locking the behaviour in: .tmp.* and other dotfiles skipped, plain and _compacted .parquet accepted, non-parquet rejected.


Note on the review loop: that's 7 rounds. The substantive findings are done — rounds 3 and 5 caught two real bugs (orphaned refs on full deletion, and silent table corruption when a column's type changes across a restart), both fixed and binary-verified. Round 7 is a one-line cleanup, which is a good signal we've converged.

The only thing that keeps recurring is WithOlderThan(0), re-flagged in rounds 1 and 6 and declined both times with the iceberg-go source and a live binary run showing retain_snapshots = 3 → exactly 3 snapshots retained out of 5. Since each review re-reads the diff fresh, I expect that one to resurface indefinitely rather than converge.

Planning to merge to the 26.09.1 collection branch unless the next pass turns up something substantive.

@xe-nvdk
xe-nvdk merged commit 249e5c6 into 26.09.1 Jul 14, 2026
xe-nvdk added a commit that referenced this pull request Jul 30, 2026
feat(iceberg): Apache Iceberg export (zero-copy, opt-in) (#532)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant