Stream PostgreSQL changes into Apache Iceberg. Written in Rust.
Status: beta. See support and status.
Documentation · Quickstart · Benchmarks · Contributing · Issues
Flow copies existing PostgreSQL rows, then continuously replicates inserts, updates and deletes into standard Iceberg tables. Query the tables directly through a compatible Iceberg reader, using your catalog and object storage.
- Postgres to Iceberg in one service. Initial COPY and ongoing logical replication, with background compaction in the same Rust process.
- Standard tables. Parquet data with Iceberg v2 position deletes or v3 deletion vectors, published through an Iceberg REST catalog to S3-compatible storage.
- Mutable data. Inserts, updates, deletes and primary-key changes, plus append-only ingestion for tables without a primary key.
- Durable recovery. A transaction journal, persistent row index and publication ledger support replay, reconnects and interrupted-bootstrap recovery.
- Maintenance alongside ingestion. Local data/delete compaction and reconciliation of external compactor rewrites.
Run the local demo with Git, Docker and Docker Compose. It includes PostgreSQL, MinIO, an Iceberg REST catalog, Trino and Flow:
git clone https://github.com/EmbrasureAI/flow.git
cd flow
docker compose -f demo/compose.yaml up --build -dOnce initialize has completed and flow has started, query the replicated table:
docker compose -f demo/compose.yaml exec trino trino \
--execute 'SELECT * FROM lake.replicated.orders'The demo guide walks through changing source rows, checking service status and cleaning up. The stack uses named volumes and does not publish ports on the host.
Install the pinned Rust toolchain, a C++ compiler, libclang, CMake, pkg-config and OpenSSL development headers, then build and validate the example configuration:
cargo build --locked --release -p flow-daemon
./target/release/embrasure-flow --config examples/flow.toml checkOn GNU/Linux, add --features jemalloc to enable process-wide allocation and
background reclamation of unused pages, including RocksDB's C++ allocations.
For a container image, pass --build-arg FLOW_FEATURES=jemalloc to docker build.
The default build uses the system allocator; the feature has no effect on other
targets. See allocator metrics for
measurement details and limits.
The Dockerfile builds a minimal image from this checkout:
docker build -t embrasure-flow .
docker run --rm --env-file flow.env \
-v "$PWD/flow.toml:/etc/flow.toml:ro" -v flow-state:/data \
embrasure-flow --config /etc/flow.toml init
docker run --env-file flow.env \
-v "$PWD/flow.toml:/etc/flow.toml:ro" -v flow-state:/data \
embrasure-flowThe image runs embrasure-flow --config /etc/flow.toml run by default as an
unprivileged user. Mount the configuration at /etc/flow.toml, set
state_dir = "/data" in it and mount a persistent volume there, and pass the
environment variables your configuration names for credentials (here through
flow.env).
Follow Getting started. In short:
embrasure-flow --config flow.toml check --source # read-only PostgreSQL preflight
embrasure-flow --config flow.toml discover >> flow.toml # generate [[tables]]
embrasure-flow --config flow.toml init # slot + initial copy
embrasure-flow --config flow.toml run # stream changesFlow needs persistent disk for its transaction journal and RocksDB row index. Readers access Iceberg independently of the running service.
flowchart LR
Postgres[PostgreSQL] -->|COPY + logical replication| Flow
Flow -->|Parquet data + deletes| Iceberg[Apache Iceberg]
Iceberg --> Readers[Compatible query engines]
Flow journals source transactions, materializes data and deletes, and commits them to the Iceberg catalog. A transaction becomes visible atomically within each destination table. A source-wide ledger advances acknowledgement only through completed transactions.
Compaction runs alongside ingestion. Before publishing a rewrite, the coordinator validates its inputs and translates intervening deletes. External physical rewrites are reconciled before ingestion reuses row locations. See the architecture and compaction protocol for the durable state transitions.
Flow is beta software. Every change runs service tests against real PostgreSQL, MinIO and an Iceberg REST catalog, reading results with stock DuckDB. On each of PostgreSQL 14–18 they cover snapshot/CDC type compatibility, schema changes, SIGKILL and catalog/object-store/PostgreSQL outage recovery, and the fail-closed checks for a lost or rewound slot or journal. Table isolation, REPLICA IDENTITY DEFAULT, lifecycle and compaction scenarios run on PostgreSQL 14 and 18; Iceberg v3, a randomized crash loop and the Trino and Spark external-maintenance checks run on PostgreSQL 18, with a longer crash loop and a resource soak nightly. Power loss (unsynced page cache) is not simulated. See the fault suite. Performance qualification is ongoing; the benchmark report records measured throughput, latency and the targets still open. Configuration and on-disk state may change between minor releases, with the upgrade path stated in the changelog; see upgrading.
Supported today:
- PostgreSQL 14–18 sources, with inserts, updates, deletes and primary-key
changes. Mutable tables need a primary key and either
REPLICA IDENTITY FULLor, when every replicated column is fixed-width,DEFAULT; see replica identity. Tables without a primary key replicate append-only. - Unpartitioned Iceberg v2 (position deletes) and v3 (deletion vectors) tables through an Iceberg REST catalog on S3-compatible storage.
- The type mappings; JSON and JSONB replicate as normalized JSON text.
- Automatic schema evolution for nullable column additions and compatible
required-to-nullable changes.
TRUNCATEand other incompatible changes block the affected table while healthy tables continue; see table isolation and operations.
Not yet supported: partitioned source tables or Iceberg partition specs, per-table re-snapshots without resynchronizing the source, cross-table query atomicity, high availability, distributed compaction and Z-order compaction. Check v3 reader compatibility and the operating limits before deploying.
| Guide | What you will find |
|---|---|
| Getting started | Build, configure, initialize and run Flow |
| Configuration | Source, storage, catalog and compaction settings |
| PostgreSQL type mappings | Supported application types across snapshot and CDC |
| Operations | Preflight, resynchronization, adding tables and planned maintenance |
| Observability | Status, HTTP probes, watermarks, metrics and diagnosis |
| Upgrading | Release compatibility and the upgrade procedure |
| Iceberg v3 | Deletion vectors, upgrades and reader compatibility |
| Code guide | Crate responsibilities and module layout |
| Integration tests | Service fixtures, reader checks and recovery scenarios |
Browse the documentation index for the full set of guides.
Bug reports, documentation improvements and code contributions are welcome. Read Contributing for development setup, checks and review expectations. Open an issue to discuss substantial changes or report a bug with reproduction steps. Participation is governed by the code of conduct. Report security issues privately as described in SECURITY.md, not in public issues.
Flow is licensed under Apache-2.0. Vendored dependencies retain their upstream licenses and attribution notices. See distribution notices for the binary license bundle.