Skip to content

About

Idempotent PostgreSQL CDC into Databricks: Rust Lambda reads a logical replication slot, Zerobus writes to Delta, Lakeflow Declarative Pipelines merges with AUTO CDC

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

1 Commit

Folders and files

Repository files navigation

Postgres CDC → Zerobus → Databricks (SDP AUTO CDC)

Change data capture from PostgreSQL into Delta Lake, extracted by a Rust AWS Lambda, written through Zerobus, then merged into a target table by a Lakeflow Declarative Pipeline (SDP) using AUTO CDC.

The point of this repo is the correctness argument: the merge is idempotent by construction, and it behaves identically whether the pipeline runs in batch (triggered) or continuous mode. Both claims are backed by tests you can run yourself.

architecture

How it works

RDS PostgreSQL  ──logical replication slot (wal2json)
      │
      │   ① peek changes   ② write to Zerobus   ③ advance the slot
      ▼
Rust Lambda (arm64) ──gRPC──▶ Zerobus ──▶ cdc_events   (Delta, append-only change log)
                                              │
                                              │ readStream
                                              ▼
                                    SDP pipeline ── AUTO CDC ──▶ customers  (SCD1)
                                       KEYS (customer_id)
                                       SEQUENCE BY lsn_num
                                       APPLY AS DELETE WHEN op = 'D'

The two design decisions that matter

1. The replication slot is the high watermark

No DynamoDB table, no parameter-store counter, no reading max(lsn) back out of Delta. PostgreSQL durably tracks confirmed_flush_lsn per slot and refuses to discard WAL you have not confirmed.

That buys three things a naive WHERE updated_at > :watermark poll cannot give you:

  • Deletes. A polling query can never see a row that is already gone.
  • A true total order. LSN is a 64-bit WAL byte offset, monotonic across the whole cluster. A wall-clock column has skew and ties.
  • No second source of truth. One watermark, owned by the database, cannot drift.

2. Read the slot with SQL functions, not the streaming replication protocol

pg_logical_slot_peek_changes()  →  write to Zerobus  →  pg_replication_slot_advance()

Two reasons.

It is what published Rust crates can actually do. tokio-postgres 0.7.18 exposes copy_in and copy_out but has no CopyBoth support, and there is no postgres-replication crate on crates.io. Streaming replication would mean depending on unreleased code or hand-rolling the wire protocol.

It gives better failure semantics. peek does not consume. So if the Zerobus write fails, the slot never advances and nothing is lost. The three steps are explicitly ordered, and the watermark moves last.

The slot is only ever advanced to the last COMMIT boundary seen, never into the middle of a transaction.

Why it is idempotent

  • Zerobus is at-least-once, so the landing table can contain duplicates. That is fine, and designed for.
  • Every change carries its LSN, unique and monotonic.
  • AUTO CDC keeps the highest sequence_by per key:
    • a replayed duplicate has an identical LSN and an identical payload, so whichever copy wins the tie the result is the same;
    • a late or out-of-order event has a lower LSN and is discarded.
  • The landing table is written by Zerobus, not by the pipeline. So a pipeline full refresh replays the entire log and converges to the same state. It cannot destroy the log it reads from.

Therefore batch vs continuous changes latency only, never correctness — ordering lives in the data, not in arrival order.

This was measured, not assumed

With the landing table holding 5,313 rows for only 204 distinct changes (heavy at-least-once duplication, including a genuine replay caused by a write that succeeded before the watermark advanced), the target still matched PostgreSQL exactly:

postgres : 5090 ba6b9cd88497740a48da011c6cb7e275 11042216420392
delta    : 5090 ba6b9cd88497740a48da011c6cb7e275 11042216420392
PARITY: PASS

Those three numbers are row count, an order-sensitive md5 over every row, and an order-independent sum of per-row md5 prefixes, computed identically on both sides. See tests/parity.sh.

Test results

30 tests executed, 30 passed, no scenario lost data. Full matrix in TEST-PLAN.md. The ones that matter:

Test Result
Deliberate double-write (170 records for 85 events) parity held
Real at-least-once replay (write succeeded, watermark did not advance) re-delivered, parity held
Late event injected with lsn_num=1 and a poison payload ignored, correct value survived
Full refresh of the target only log replayed from scratch, parity held
Invalid Zerobus credentials 401, watermark did not move, exit code 1
Recovery after that failure all at-risk changes delivered
Continuous mode converged ~45 s, no manual trigger
Batch ↔ continuous switch, no full refresh parity held both directions

Reproduce the interesting ones:

bash scripts/generate-changes.sh 50 30 5                               # 50 insert / 30 update / 5 delete
bash scripts/run-local.sh '{"mode":"stream","duplicate_writes":true}'  # force an at-least-once replay
bash scripts/run-local.sh '{"mode":"stream","skip_advance":true}'      # simulate a crash before the watermark commits
bash tests/parity.sh                                                   # must still pass after either

duplicate_writes and skip_advance are deliberate test switches in the extractor, not production behaviour.

Prerequisites

  • An AWS account where you can create RDS, VPC resources and Lambda
  • A Databricks workspace on AWS with Unity Catalog and a serverless SQL warehouse
  • aws CLI, databricks CLI, psql, Rust 1.70+, and cargo-lambda:
    brew install cargo-lambda/tap/cargo-lambda zig

Put RDS in the same region as your Databricks workspace so Zerobus traffic stays in-region.

Reproducing it

# 0. configure
cp infra/config.env.example infra/config.env
$EDITOR infra/config.env          # fill in your AWS profile, workspace, catalog, warehouse

# 1. Postgres source + networking
bash infra/01-create-rds.sh          # ~8 min; generates the master password into
                                     # infra/config.local.env (gitignored)
bash infra/02-create-vpc-egress.sh   # ~2 min; private subnet + NAT gateway + Lambda SG

# 2. allow your own IP to reach Postgres, then load the source schema
source infra/config.env; source infra/config.local.env
MYIP=$(curl -s https://checkip.amazonaws.com)
SG=$(aws ec2 describe-security-groups --filters "Name=group-name,Values=$PG_SG_NAME" \
       --query 'SecurityGroups[0].GroupId' --output text)
aws ec2 authorize-security-group-ingress --group-id "$SG" \
  --protocol tcp --port 5432 --cidr "$MYIP/32"

export PGPASSWORD="$PG_MASTER_PASSWORD"
PG="host=$PG_HOST port=$PG_PORT dbname=$PG_DBNAME user=$PG_MASTER_USER sslmode=require"
psql "$PG" -f sql/01-postgres-setup.sql
psql "$PG" -c "GRANT rds_replication TO $PG_MASTER_USER;"    # required: see Gotchas

# 3. Databricks landing table
bash scripts/apply-databricks-setup.sh

# 4. a service principal for Zerobus, then grant it write on the landing table
#    (client id + secret go into infra/config.local.env as DBX_CLIENT_ID / DBX_CLIENT_SECRET)
#    GRANT USE CATALOG / USE SCHEMA / SELECT, MODIFY to that principal.

# 5. build and run the extractor
cd lambda && cargo test && cargo lambda build --release --arm64 && cd ..
bash scripts/run-local.sh '{"mode":"snapshot"}'   # initial load of existing rows
bash scripts/run-local.sh '{"mode":"stream"}'     # drain the change log

# 6. the merge
cd pipeline && databricks bundle deploy -t batch -p "$DBX_PROFILE"
databricks bundle run pgcdc_merge -t batch -p "$DBX_PROFILE" && cd ..

# 7. prove it
bash tests/parity.sh

The same binary runs as a Lambda or as a local CLI: it checks for AWS_LAMBDA_RUNTIME_API and falls back to one-shot CLI mode, so you can iterate without deploying. To deploy it:

bash infra/04-deploy-lambda.sh <lambda-execution-role-arn>
bash infra/05-schedule.sh

The role must trust lambda.amazonaws.com and allow ec2:CreateNetworkInterface, ec2:DescribeNetworkInterfaces, ec2:DeleteNetworkInterface (i.e. the AWS managed policy AWSLambdaVPCAccessExecutionRole) plus CloudWatch Logs.

Batch or continuous

One line in pipeline/databricks.yml: continuous: true plus pipelines.trigger.interval. Correctness is unchanged by design; see above.

Note that on an EventBridge schedule the freshness floor is 1 minute (its minimum rate), not Zerobus, which is P50 ≤5 s time-to-table. Genuinely sub-minute freshness needs a long-running extractor rather than a scheduled Lambda.

Gotchas found the hard way

  • The landing table is a streaming source, so treat it as immutable. Zerobus only appends, which is what makes this work. But any manual DML on cdc_events creates a non-append commit and the streaming read then fails fatally with DELTA_SOURCE_IGNORE_DELETE, unable to restart without a full refresh. I hit this by "tidying up" a single injected test row with a DELETE. The pipeline now reads WITH (SKIPCHANGECOMMITS) as a guard, but the operational rule stands: never DML the change log. Recovery is a full refresh, which is safe precisely because the log lives outside the pipeline.
  • pg_available_extensions does not list wal2json. It is an output plugin with no SQL objects, so that view is the wrong check. Test it by creating a slot. wal2json 2.6 is available on RDS PostgreSQL 16.
  • The RDS master user is NOREPLICATION by default. Run GRANT rds_replication TO <user> before creating a slot.
  • Binding a Rust String to a pg_lsn parameter fails. Cast in SQL: $2::text::pg_lsn.
  • RDS presents a certificate from a private Amazon CA. The public webpki roots will not verify it, so the published RDS CA bundle is embedded (lambda/src/rds-global-bundle.pem) and verified properly, rather than disabling verification.
  • CREATE TEMPORARY STREAMING VIEW is not valid SDP SQL. Use CREATE TEMPORARY VIEW and put STREAM() in the FROM clause.
  • The two ingest paths emit different timestamp formats. The snapshot path (to_jsonb) produces 2026-09-16T13:14:39.33375+00:00; wal2json produces 2026-09-17 13:17:27.868965+00. The pipeline normalises both to ISO-8601 instead of guessing a single format string.
  • Delta TIMESTAMP over Zerobus is int64 epoch MICROSECONDS, not milliseconds.
  • sequence_by cannot be a STRUCT. Use one orderable scalar; here a BIGINT LSN.
  • Do not set delta.autoOptimize.optimizeWrite on serverless — it raises DELTA_UNKNOWN_CONFIGURATION. Serverless does it for you.
  • An unconsumed replication slot retains WAL and can fill your RDS storage. Cap it with max_slot_wal_keep_size (dynamic on RDS, no reboot) and understand the trade-off: past the cap PostgreSQL invalidates the slot, so you lose changes rather than the instance. Size it to your worst-case extractor outage and alarm on CloudWatch OldestLogicalReplicationSlotLag / ReplicationSlotDiskUsage well before it.
  • A VPC Lambda in a public subnet has no internet egress (its ENI gets no public IP), and VPC endpoints only front AWS services, not an arbitrary public gRPC endpoint. Hence the NAT gateway.
  • format_string('%.2f', <decimal>) fails in Spark with f != Decimal; cast the decimal to string instead.
  • Spark's array_sort on strings sorts lexicographically. A cross-system sorted-hash comparison gives a false mismatch unless you zero-pad the key.

Table properties, and why

The merge target uses:

CLUSTER BY (customer_id)
TBLPROPERTIES (
  'delta.enableDeletionVectors' = 'true',
  'delta.enableRowTracking'     = 'true',
  'delta.targetFileSize'        = '67108864'   -- 64 MB
)
  • Cluster on the merge key. The MERGE always matches on the primary key, and a BIGINT key clusters into narrow min/max ranges so file pruning actually works. A UUIDv4 key would give roughly zero pruning, because every file's range spans the whole keyspace.
  • Deletion vectors give low-shuffle merge: mark rows deleted in a sidecar instead of rewriting whole files for single-row updates, which is the dominant CDC pattern.
  • Smaller target file size than the 1 GB default, so each rewrite touches fewer rows.

Cost

Resource Approximate
RDS db.t4g.micro, 20 GB gp3 ~$15/month
NAT gateway ~$32/month + $0.045/GB
Lambda (1/min, 512 MB, ~2 s) pennies
SDP serverless triggered: per update. Continuous: always-on

infra/99-teardown.sh removes the networking and compute, which stops the NAT charge. It deliberately stops short of deleting the database and prints that step for you to run yourself.

Layout

Path What
infra/01-create-rds.sh RDS PostgreSQL + parameter group (logical replication) + SG
infra/02-create-vpc-egress.sh Private subnet + NAT gateway + Lambda SG
infra/03-secrets-and-iam.sh Parameter-store secrets + Lambda role
infra/04-deploy-lambda.sh Package and deploy the arm64 Lambda into the VPC
infra/05-schedule.sh EventBridge rule
infra/99-teardown.sh Remove networking/compute
sql/01-postgres-setup.sql Source schema, table, 5,000 synthetic rows
sql/02-databricks-setup.sql Landing table (the Zerobus sink)
lambda/src/pg_cdc.rs Slot management, peek, advance, snapshot
lambda/src/wal2json.rs wal2json v2 decode, LSN ↔ i64 (unit tested)
lambda/src/zerobus_sink.rs Zerobus JSON ingest
pipeline/src/cdc_merge.sql The AUTO CDC merge
tests/parity.sh Acceptance oracle: Postgres vs Delta fingerprint
docs/zerobus-vs-lakeflow-connect.md When to use this vs the managed connector
docs/architecture-prompt.txt Regenerates the architecture diagram

Should you build this?

Often, no. Databricks ships a managed PostgreSQL connector in Lakeflow Connect that requires no extractor code at all, and for most workloads that is the right answer.

This pattern earns its keep when you need materially lower latency than a triggered 5-minute schedule, or you need to own the extraction layer (custom masking, routing, filtering at source). It costs you an extractor to operate.

docs/zerobus-vs-lakeflow-connect.md is an honest head-to-head, including the cases where the managed connector wins.

Status

The extraction and merge path is built and verified end to end against real RDS and real Zerobus. Everything in TEST-PLAN.md marked with a result was actually run.

The Lambda deployment scripts are written and validate their inputs, but were not executed in the environment this was developed in, because creating the execution role required IAM permissions that were not available there. Treat infra/04-deploy-lambda.sh and infra/05-schedule.sh as untested in anger.

Licence

MIT. See LICENSE.

About

Idempotent PostgreSQL CDC into Databricks: Rust Lambda reads a logical replication slot, Zerobus writes to Delta, Lakeflow Declarative Pipelines merges with AUTO CDC

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages