Blog

GCP Kafka to Iceberg: Retention, Table Semantics, and Replay Boundaries

Table of Contents

Table of Contents

An Apache Kafka topic and an Apache Iceberg table can contain the same business events while making different promises. Kafka gives producers and consumers an ordered log, offsets, and retention rules. Iceberg gives analytical readers snapshots, manifests, schema evolution, and table commits. A GCP pipeline that connects the two is reliable only when those contracts are explicit.

The practical question is not “How do I copy Kafka into Iceberg?” It is “Which event is durable where, for how long, and what can a replay prove?” A connector or native writer may successfully append files while a consumer still expects Kafka offsets, a table reader expects snapshot isolation, and a retention job deletes the source before a backfill is complete. The design needs a boundary map before it needs a tool choice.

Kafka records moving through a GCP stream-to-table path into an Iceberg table

1Start with two contracts, not one pipeline

The Kafka contract is about the log. Producers choose a topic and partition key; brokers assign offsets; consumers commit positions; retention determines when records stop being available. Ordering is normally meaningful within a partition, while a consumer group’s progress is represented by offsets rather than table snapshots. A record can be visible to one consumer and still be waiting for a table writer to process it.

The Iceberg contract is about a table state. A writer creates data files and manifests, then publishes a snapshot through a catalog. Readers see a committed table version according to the table format and engine they use. A snapshot is not a Kafka offset, and a table commit is not automatically a confirmation that every downstream consumer has processed the record.

Write these boundaries into the runbook:

  • Record identity: Which fields identify an event across retries, duplicate delivery, and compaction?
  • Progress marker: Is recovery based on a Kafka offset, a connector checkpoint, an Iceberg snapshot ID, or a combination?
  • Visibility point: When does a record become readable from Kafka, from the table, and from both?
  • Retention owner: Which policy removes the record, its data file, its manifest, or its catalog metadata?

The answers prevent a common failure mode: treating “the table has data” as proof that the Kafka-to-Iceberg handoff is complete.

2Map the GCP data path and its storage boundaries

On Google Cloud, draw the path as separate services even when one product hides the implementation. A typical design has Kafka producers and consumers in a VPC, a Kafka service retaining the log, a writer reading records, a catalog publishing Iceberg snapshots, and Cloud Storage holding data and metadata files. Query engines read table snapshots from the catalog and object storage. Network paths, IAM principals, and failure domains are different at each edge.

Google’s Managed Service for Apache Kafka overview describes the managed Kafka boundary. It does not, by itself, define an Iceberg table contract or a connector’s replay behavior. The table layer should be reviewed against the Apache Iceberg specification and the selected writer and catalog documentation. Cloud Storage’s consistency model is also relevant, but storage consistency does not remove the need to reason about catalog commit ordering.

Keep the following assets separate in the ownership matrix:

BoundaryQuestion to answerEvidence to keep
Kafka logHow long can a consumer replay the source records?Topic retention, partition offsets, consumer checkpoint
WriterWhat happens after a timeout or duplicate delivery?Retry policy, idempotency key, dead-letter handling
Iceberg tableWhich snapshot is complete and queryable?Snapshot ID, manifest list, commit log
Cloud StorageWho can read, write, delete, or restore files?Bucket policy, service account, object versioning policy
CatalogWhere is the authoritative table pointer?Catalog endpoint, namespace, commit audit

This inventory makes a GCP outage discussion concrete. A Kafka broker restart, a writer crash after file upload, and a catalog commit conflict are different incidents even if all appear as “the pipeline is delayed.”

3Retention is a replay budget

Kafka retention and table retention answer different questions. Kafka retention limits how far a consumer can seek back in the log. Iceberg retention controls how long snapshots and their referenced files remain available for time-travel queries or table repair. A table may retain a snapshot after Kafka has deleted the source record, or Kafka may retain records that have not yet been published in a table snapshot.

Define at least three horizons for each workload:

  1. Operational replay: the period in which the writer can re-read Kafka records after a transient failure.
  2. Analytical history: the period in which a query engine can read an Iceberg snapshot and its files.
  3. Rebuild window: the period in which the team can reconstruct a table after a bad transform, schema mistake, or accidental delete.

These horizons do not need the same duration. They do need an explicit handoff. If the writer is paused longer than Kafka retention, a table backfill cannot depend on the original topic unless another source exists. If old Iceberg files are expired while a long-running query still references the snapshot, the table maintenance policy may break the query’s expected time-travel window.

Use a retention worksheet with the event timestamp, Kafka offset range, table snapshot ID, and file locations for one representative batch. Record what happens when a batch is retried, partially committed, or replayed after a schema change. That evidence is more useful than a single “retention: 7 days” field because it shows where the data actually remains available.

4Table semantics decide whether replay is safe

A stream of records is not automatically a table of rows. The conversion has to define how keys, deletes, updates, late events, and duplicate deliveries map to table operations. An append-only topic can often map to append-only data files. A change-data-capture stream may require an upsert key, delete markers, or merge logic. A table that looks correct in a sample query can still be wrong when events arrive out of order.

Ask these questions before selecting a sink:

  • Is the table append-only, or can a later event replace an earlier row?
  • Which Kafka fields survive conversion: key, value, headers, partition, offset, and event time?
  • Does the writer preserve nullability, logical types, and schema version metadata?
  • How are deletes represented, and what does a replay do with an already-applied delete?
  • Is deduplication keyed by Kafka offset, event ID, or a business key?

The Iceberg Kafka Connect documentation is a useful reference for one connector-based path, but its configuration does not define the semantics of every writer or transformation. Treat the selected implementation as a contract to test. The acceptance test should compare source records with table rows, table snapshots, and duplicate behavior, rather than only file counts.

Replay flow showing Kafka offsets, table snapshots, and the commit boundary

5Test replay across the commit boundary

A reliable replay test injects failures at the awkward moments, not only after a clean shutdown. Stop the writer after it reads records but before it commits a table snapshot. Stop it after files are uploaded but before the catalog pointer moves. Force a retry after a commit response times out. Then restart from the recorded checkpoint and inspect both Kafka and Iceberg outcomes.

The expected result depends on the writer’s guarantees. A retry may produce duplicate rows in an append-only table, or it may be safe when the writer uses an idempotent event key. A commit conflict may require a fresh snapshot read and a new commit attempt. A timeout does not tell the operator whether the commit succeeded; the catalog and table metadata are the authority.

Keep replay evidence in a small fixture:

  • Source topic, partition, offset range, and event IDs
  • Writer checkpoint before the failure
  • Uploaded file and manifest identifiers, if visible
  • Iceberg snapshot ID before and after recovery
  • Row-level result for append, update, and delete cases
  • Query output from a reader started before and after the replay

A replay is complete when the team can explain every source offset and table row in the fixture. “The connector is running again” is a service state, not data proof.

6Where AutoMQ can fit and what to verify

Once the contracts are clear, a Kafka-compatible stream platform with a native stream-to-table option becomes one architecture to evaluate. AutoMQ describes Table Topic as a built-in capability that materializes Kafka topics as Apache Iceberg tables on object storage. Its Table Topic overview provides the current product boundary and configuration context.

That option changes the number of moving parts, but it does not remove the need for a contract. Confirm the supported AutoMQ deployment mode and GCP storage integration for your target release, then test retention, schema mapping, snapshot publication, replay, and access control with your own data. If a stream needs joins, stateful transformations, or event-time computation, a stream processor may still be the right boundary; a native table path is best evaluated for data that can be materialized directly.

The useful comparison is therefore operational:

Readiness matrix for Kafka-to-Iceberg retention, semantics, replay, and ownership

Decision areaConnector or stream processorNative stream-to-table path
TransformationFlexible code and stateful processingBest for direct materialization with defined mappings
Failure surfaceKafka, runtime, checkpoint store, catalog, and storageKafka-compatible runtime, table writer, catalog, and storage
Replay proofCheckpoint plus source offset plus table outcomeSource offset plus table snapshot and row outcome
OwnershipSeparate runtime and stream platform teamsFewer runtime boundaries, but product-specific checks

Treat this table as a review frame, not a feature score. The correct choice depends on the semantic work the pipeline must perform and the evidence the lakehouse team requires.

7A release gate for GCP Kafka to Iceberg

Before production traffic flows into an Iceberg table, require a signed readiness packet:

  1. Source contract: topic, partitions, keys, event-time fields, and Kafka retention are documented.
  2. Table contract: schema, identifiers, update/delete behavior, partitioning, and snapshot policy are documented.
  3. Identity contract: deduplication and retry behavior are proven with duplicate and timeout tests.
  4. Replay contract: an offset range can be replayed into a known table snapshot without unexplained rows.
  5. Access contract: GCP IAM, bucket permissions, catalog ownership, and audit logs identify every writer and reader.
  6. Expiry contract: Kafka retention, snapshot expiration, and object cleanup cannot delete data before the rebuild window closes.

The final approver should be able to answer one question: if a table query is wrong tomorrow, can the team identify the source offsets, reconstruct the intended snapshot, and explain which retention policy still preserves the evidence?

7.1FAQ

7.1.1Does Kafka retention guarantee Iceberg history?

No. Kafka retention governs the source log; Iceberg retention governs snapshots and files. Coordinate the two only after defining the replay and rebuild windows.

7.1.2Is an Iceberg snapshot equivalent to a Kafka offset?

No. An offset identifies a position within one Kafka partition. An Iceberg snapshot identifies a committed table state. A replay design may need both markers.

No. Connector-based or native materialization can fit direct mappings. Stateful joins, event-time logic, and complex transformations still require a processing layer suited to that work.

7.1.4Can I claim GCP support from a generic Iceberg integration page?

No. Verify the target deployment mode, GCP region, storage path, catalog, IAM, and current product release. Keep provider capability claims tied to the specific evidence you tested.

A GCP Kafka-to-Iceberg design is ready when its handoff can be replayed and explained. If your team wants to test a Kafka-compatible stream-to-table architecture against these gates, start with the AutoMQ evaluation path and bring the retention worksheet, replay fixture, and ownership matrix to the review.

Newsletter

Subscribe for the latest on cloud-native streaming data infrastructure, product launches, technical insights, and efficiency optimizations from the AutoMQ team.

Join developers worldwide who leverage AutoMQ's Apache 2.0 licensed platform to simplify streaming data infra. No spam, just actionable content.

I'm not a robot
reCAPTCHA

Never submit confidential or sensitive data (API keys, passwords, credit card numbers, or personal identification information) through this form.