Change-data-capture (CDC) is the backbone of modern analytics: it keeps your lakehouse fresh, powers near-real-time analytics and feeds downstream ML and BI systems. This guide walks data engineers and analytics engineers through a concrete, production-ready pattern for reliable CDC ingestion from OLTP databases (MySQL/Postgres) into Apache Iceberg using Debezium, Apache Kafka, and Apache Flink. It focuses on upsert semantics, schema evolution, tombstones/deletes, compaction, operational settings, testing and common pitfalls.

Why this pattern?

This stack is popular because each component plays a clear role:

  • Debezium: low-latency connectors that capture row-level changes from databases and publish them to Kafka.
  • Kafka: durable, replicated event backbone and buffer for replays and multi-consumer topologies.
  • Flink: stream processor that materializes CDC into an analytical table format — Apache Iceberg — handling upserts and schema changes with low latency.

The goal is a robust pipeline that produces correct, compact analytical tables in Iceberg suitable for ad-hoc SQL, BI and ML.

Architecture overview

At a high level:

  1. Debezium connectors tail database WALs and publish change events (keyed by primary key) to compacted Kafka topics.
  2. Kafka acts as a replayable buffer; messages include schema metadata via Avro/Protobuf (Schema Registry) or Debezium JSON.
  3. Flink reads topics as a changelog stream, applies idempotent upserts or MERGE operations against Iceberg tables using Flink’s Table API/SQL + Iceberg connector.
  4. Iceberg handles file-level metadata, compaction/optimization and time travel for analytics.

Prerequisites & versions

  • Debezium connector for your source DB (MySQL/Postgres): deployment in Kafka Connect (recommended) or Debezium engine embedded.
  • Apache Kafka cluster with topic compaction enabled for CDC topics, and Schema Registry (Confluent/Apicurio) if using Avro/Protobuf.
  • Apache Flink cluster supporting the Iceberg connector and two-phase commit sinks (Flink 1.16+ or recent stable release in 2026).
  • Apache Iceberg table storage on object store (S3/GCS/Azure) or HDFS, plus metadata catalog (Hive Metastore, AWS Glue, or Iceberg-native catalog).
  • Production monitoring (Prometheus/Grafana), alerting and capacity for Flink state/backups.

Step 1 — Debezium: capture and produce Kafka-friendly records

Key decisions at the Debezium layer:

  • Message key must be the table primary key(s). Kafka topic compaction works only when keys are stable.
  • Choose value format: Avro/Protobuf with Schema Registry or Debezium JSON. Avro/Protobuf is preferred for compactness and strict schema evolution.
  • Set snapshot.mode appropriately (for initial load): commonly "initial" to capture baseline state, then stream changes.
  • SMTs (Single Message Transforms): use the "unwrap" transform to emit a straightforward after/before payload structure; emit tombstones on deletes (so compaction removes deleted keys).

Operational tips:

  • Use topic per table (default) so Flink can map topics to Iceberg tables cleanly.
  • Enable Kafka topic compaction; set retention to -1 or appropriate for your retention model.
  • Tune connector tasks based on DB scale; avoid large snapshot loads during peak hours.

Step 2 — Topic and schema hygiene

Recommended practices:

  • Store value schemas in a central registry; evolve using backward/forward-compatible rules. Schema evolution in Avro/Protobuf avoids downstream decode issues.
  • Partition keys in Kafka: use default partitioning unless you need ordering by key—Flint will handle parallelism.
  • Produce tombstone messages (null value) for deletes so Kafka compaction removes keys; this is key when turning changelog into final table state.

Step 3 — Flink: read changelog and write to Iceberg as upserts

Flink's role is to convert the event stream into table updates in Iceberg with transactional guarantees.

Design options:

  • Stream-to-Upsert: define Kafka source as a changelog stream, then execute SQL MERGE INTO or use Iceberg’s upsert capabilities to apply updates/deletes.
  • Table per Kafka topic: create a Flink table that maps the topic schema (including key columns) and treat the stream as a changelog with INSERT/UPDATE_BEFORE/AFTER and DELETE events.

Key Flink settings and behaviors:

  • Enable checkpointing and two-phase commit (2PC) support in the Iceberg sink and Kafka source to achieve end-to-end exactly-once.
  • Set state backend (RocksDB recommended) and checkpoint interval based on throughput and latency requirements.
  • Use Flink’s built-in Debezium/Avro formats or convert Avro records into a relational schema. Map primary key columns explicitly.
  • Implement an idempotent upsert strategy: Flink + Iceberg combination should ensure that repeated application of the same message does not corrupt state.

Applying upserts

Two pragmatic approaches to apply CDC changes to Iceberg from Flink:

  • Streaming MERGE: Use Flink SQL to run continuous MERGE INTO operations when supported by your Flink/Iceberg integration. This maps inserts/updates/deletes to Iceberg table updates.
  • Changelog write: Use Iceberg’s changelog writer sink which accepts a changelog stream (INSERT/UPDATE/DELETE). Iceberg will apply changes directly when the sink supports transactional semantics.

Pick the method that matches your connector capabilities; many teams find the changelog writer simpler and faster for high-throughput workloads.

Step 4 — Iceberg table design and compaction

Iceberg is optimized for analytics, but CDC workloads have specific needs:

  • Primary key mapping: Iceberg does not enforce primary keys at the file level. Use partitioning and an explicit primary-key-like column set in your application logic, and ensure Flink/Merge applies changes deterministically.
  • File sizing and small-file control: high-frequency updates can create many small files. Configure writer target file size and use periodic compaction/optimize jobs to rewrite small files into larger ones.
  • Compaction strategy: run scheduled rewrite-manifest/major-compaction during low-load windows. Iceberg’s data rewriting utilities or cloud-optimized compaction jobs work well.
  • Retention & vacuuming: tombstones and old snapshots accumulate. Define snapshot expiration policy and run vacuum jobs to free storage.

Step 5 — Exactly-once, transactions and recovery

To ensure correctness across Debezium → Kafka → Flink → Iceberg:

  • Debezium → Kafka relies on the connector’s at-least-once semantics; Kafka provides durable storage and compaction.
  • Flink must be configured for checkpointing and sinks must implement 2PC or transactional writes. The Iceberg sink integrates with Flink’s 2PC to commit table changes only on successful checkpoints.
  • Design for idempotency: even with transactional sinks, ensure your MERGE logic or changelog application can handle replayed messages safely.

Failure scenarios to plan for: Flink job restarts, Kafka partition rebalances, Debezium connector restarts and database failovers. Regularly test end-to-end failure and recovery procedures in staging.

Step 6 — Schema evolution and incompatible changes

Schema changes are common in OLTP systems. Handle them deliberately:

  • Use schema registry to manage changes; prefer additive changes (new nullable columns) and avoid breaking renames or type narrows unless coordinated.
  • When columns are renamed, consider publishing a mapping layer in Flink (rename during ingest) or use an ETL step to reconcile names before writing to Iceberg.
  • Iceberg supports schema evolution; ensure your Flink-to-Iceberg mapping is driven by schema-aware deserializers and can handle missing/new fields.

Step 7 — Testing, rollout and backfill

Rollout plan:

  1. Set up a staging pipeline with a representative subset of tables and traffic.
  2. Run initial snapshot load via Debezium snapshot, or perform a controlled backfill into Iceberg if you already have historical data.
  3. Enable topic compaction in Kafka and validate tombstone behavior for deletes.
  4. Start Flink job in 'shadow' or read-only mode that writes to a test Iceberg table; compare materialized state to source DB via spot checks and row counts.
  5. Gradually move consumer workloads to the Iceberg tables and cut over when parity is confirmed.

Step 8 — Monitoring & troubleshooting

Key metrics and alerts:

  • Debezium: connector lags, snapshot progress, connector task failures.
  • Kafka: partition under-replicated, end-to-end latency, compaction backlog, broker CPU/I/O, consumer lag per topic/partition.
  • Flink: checkpoint success/failure, operator backpressure, state size, restart rate, sink commit failures.
  • Iceberg: number of small files, snapshot growth, metadata table size, vacuum/compaction failures.

Troubleshooting tips:

  • If you see duplicate rows: verify Kafka keying, tombstone deletes, Flink checkpointing and idempotent merge logic.
  • If updates are missing: check that Debezium is emitting deletes/tombstones and that Flink is handling UPDATE_BEFORE/UPDATE_AFTER messages correctly.
  • If small-file explosion occurs: reduce target file size in writer, increase compaction frequency, and batch low-priority updates.

Performance tuning & cost

Real-world tuning areas:

  • Adjust Flink parallelism to match Kafka’s partition count and downstream I/O capacity.
  • Use vectorized writers and optimized file formats (Parquet/ORC) with columnar compression tuned to your query patterns.
  • Avoid excessive snapshots and keep Iceberg metadata compact by regularly running metadata compaction jobs.
  • Balance data freshness with cost: micro-batching in Flink can reduce small-file churn at the expense of a few seconds of latency.

Security & governance

Security checklist:

  • Encrypt Kafka topics in transit (TLS) and at rest (disk encryption/OAuth for cloud brokers as needed).
  • Use authenticated Schema Registry and RBAC to prevent unauthorized schema evolution.
  • Control Iceberg catalog access via IAM (cloud) or Ranger/Atlas integrations; enforce least privilege for writers and readers.
  • Log provenance: track snapshot IDs, Flink job IDs and checkpoints used to commit Iceberg snapshots for audits.

Common pitfalls & how to avoid them

  • Misconfigured keys—if the Kafka key doesn’t exactly match the DB primary key, compaction and upserts will be wrong. Always verify key mapping in Debezium.
  • Ignoring tombstones—without tombstones deletes won’t compact in Kafka and will persist in Iceberg unless explicitly handled.
  • Under-provisioned Flink state—state blow-ups during heavy updates will cause frequent checkpoints and restarts. Monitor state sizes and tune RocksDB and memory.
  • Schema drift without coordination—uncoordinated renames or type changes break deserialization; enforce schema evolution rules via registry policies.

Operational checklist

  • Debezium: validate primary-key mapping, snapshot mode, and SMT unwrapping; enable tombstones.
  • Kafka: enable compaction, validate retention, register schemas in registry.
  • Flink: enable checkpointing and 2PC, choose correct Flink-Iceberg connector and state backend (RocksDB), set parallelism and checkpoint intervals.
  • Iceberg: set target file size, schedule compaction/vacuum jobs, configure snapshot retention.
  • Testing: run parity checks, failure-recovery drills, and schema-evolution tests before production cutover.

Conclusion

Debezium → Kafka → Flink → Iceberg is a battle-tested pattern for bringing OLTP changes into analytics tables. The engineering work centers on correct schema management, producing compact and keyed Kafka topics, configuring Flink for transactional commits, and keeping Iceberg healthy through compaction and retention. With careful rollout, robust monitoring and routine maintenance tasks (compaction, vacuum, schema governance), this pattern yields accurate, low-latency analytical tables ready for BI and ML workloads.

Next steps: prototype with a single, medium-change-rate table; validate end-to-end behavior for inserts, updates, deletes and schema changes; then iterate on compaction and operational playbooks.