Table of Contents
Table of Contents
A CDC pipeline can be connected and still be wrong. A source connector may read a database log, Kafka may accept the records, and a sink may report successful writes while an update arrives out of order or a retry applies the same change twice. Those failures are hard to see because every component can show a healthy status at the same time.
The contract is narrower than “the connector is running.” Identify the order the source preserves, the order Kafka guarantees, where retries can duplicate records, and how the sink proves its final state. Treat those boundaries as evidence gates before choosing a GCP connector or managed service.
1Start with the record contract, not the connector catalog
Change data capture turns database log entries into events. The Debezium architecture guide is a useful reference for the source-to-event boundary. The event should carry enough identity for a downstream system to decide whether it is seeing a first application, a retry, or an older version of the same entity. At minimum, define the source table or aggregate key, operation type, source position, transaction identifier when available, and the schema version used to decode the payload.
That contract states what “correct” means. A ledger may require all changes for one account to remain in order. A search index may tolerate an intermediate update if the latest version wins. A warehouse load may need replayable history, even when the table uses upserts. Those contracts differ, so connector availability alone cannot approve the design.
A useful contract records four decisions:
- Identity: Which key ties an event to the row or aggregate it changes? Keep the key stable across updates and deletes.
- Position: Which source log position or transaction marker lets an operator locate the event again? Treat this as evidence, not as a display-only field.
- Version: Which sequence, commit timestamp, or source version lets the sink reject an older update? A wall-clock timestamp alone may not define a total order.
- Outcome: Is the sink expected to append history, apply an idempotent upsert, or produce a side effect exactly once within its own boundary?
2Ordering exists at boundaries
Database logs usually preserve an order for changes in the log they expose, but that does not automatically become one global order in Kafka. Kafka semantics define ordering within a partition, so Kafka preserves record order there. If events for one customer or account must stay ordered, the producer or source connector needs a stable partition key that keeps that aggregate on one partition. Events for different keys may be processed concurrently, and that is usually the desired trade-off.
Transaction boundaries need their own test. A database transaction can update several rows, while a CDC connector emits several Kafka records. Consumers may observe those records one at a time unless the source, connector, and consumer contract carries transaction metadata or an equivalent grouping signal. Do not describe this as global ordering unless the complete path proves it.
The same distinction applies to retries. A producer retry can resend a record; a connector task restart can repeat a batch; a sink retry can reapply a previously acknowledged operation. Idempotent producer settings can reduce duplicate writes at the Kafka producer boundary, but they do not make an external sink idempotent. The sink still needs a key and a write rule that can recognize a repeated event.
The ordering test should therefore name the boundary under test. A compact matrix keeps the discussion concrete:
| Boundary | Question | Evidence to capture |
|---|---|---|
| Source log | Can an operator locate the exact change again? | Source position, transaction marker, operation, and schema version |
| Producer and connector | Does a retry preserve key and ordering assumptions? | Retry and restart run with duplicate and out-of-order checks |
| Kafka topic | Do related events stay in one partition? | Key-to-partition samples and offset sequence for one aggregate |
| Consumer | Does the consumer commit after the required work? | Commit policy, restart test, and offset-to-output comparison |
| Sink | Does a repeated event leave the intended final state? | Idempotency key, write result, and read-back verification |
The table separates facts often collapsed into one “delivery guarantee.” Each row has a different owner and failure signal, which lets a team stop a bad replay before it reaches a customer-facing system.
3Choose a GCP path by its recovery evidence
On Google Cloud, a CDC design may combine a source connector, Kafka Connect workers, a Kafka cluster, and a sink runtime such as Dataflow or another consumer service. Google Cloud's Datastream overview describes its change-data-capture service and supported destinations. The Kafka Connect documentation helps map the connector runtime boundary.
For a Kafka Connect path, test worker restart, task rebalancing, offset storage, plugin configuration, and credentials as one recovery sequence. For a Dataflow path, test checkpoint recovery, Kafka consumer offsets, and the destination write behavior together. If you use Google Cloud Managed Service for Apache Kafka, keep the provider's broker responsibility separate from the connector and sink responsibilities. Managed brokers can reduce cluster operations while leaving replay and data correctness in your application path.
The decision record should answer three questions:
- Where is the authoritative cursor: the source log position, Kafka offset, connector offset store, or a combination?
- Which component can pause safely when the sink is unavailable, and where is the backlog measured?
- Which state must be restored together so a replay cannot mix an old schema, a later offset, and a partially applied sink write?
A design that cannot answer those questions has a connectivity plan, not a recovery plan.
4Replay is a controlled experiment
Replay should begin with a bounded range and a known result. Record source position or Kafka offsets, schema version, sink state before the run, and the rule for repeated records. Isolate replay traffic from the live consumer group unless the sink contract supports sharing that load.
A useful replay run exercises the failure modes that are otherwise hidden in a green-path demo:
- Stop the sink after it has acknowledged some records but before the consumer commits the corresponding offset.
- Restart the connector or consumer and observe whether the same records appear again, in which order, and with which keys.
- Inject a schema change at a controlled position, then replay records on both sides of that change.
- Compare the sink's final state and audit history with a source-side query or an independently computed expected result.
The pass condition is not “the task reached running.” It is a reconciliation statement: every expected change appears, a repeated change creates no unintended side effect, and an older version cannot overwrite a later version. Capture the event trace and final sink state, because either view can hide an error.
Replay also changes capacity. A sink that handles live traffic may fall behind when it reads an older range, and a backlog can compete with fresh events for the same connector tasks or destination quota. Measure backlog age and drain behavior during the replay window. Keep those measurements as workload evidence rather than publishing a universal throughput claim.
5Sink correctness belongs to the destination owner
The sink owns the last irreversible step, so its write model must be explicit. An append-only destination can record every event, but an upsert destination must decide how to handle duplicate or late events. A sink writing to a warehouse may use a merge keyed by the source identity and version. A side-effecting service may need a durable idempotency key and a response log before it acknowledges the Kafka offset.
Do not use “exactly once” as a pipeline-wide label unless every boundary has a documented transaction or idempotency mechanism. Kafka transactions can coordinate Kafka records and offsets within Kafka's transaction model. They do not automatically include a database, HTTP API, or other external destination. The honest description is usually “at-least-once transport with an idempotent sink,” followed by the evidence that demonstrates the sink rule.
A sink review should cover failure and repair ownership:
- Duplicate event: Which field identifies the operation, and what does the sink do when it sees it again?
- Late event: Can an older source version overwrite a later one? If not, where is the comparison made?
- Partial batch: What happens when half a batch is accepted and the task fails before commit?
- Poison record: Can one malformed event be quarantined without advancing past unrelated records?
- Repair: Can an operator re-run a bounded range and verify the resulting state without manual edits?
These controls matter more than a connector feature list because they describe what operators can prove after failure.
6Where AutoMQ fits
Once the record contract, ordering boundaries, replay range, and sink checks are defined, the Kafka storage layer becomes one part of the evaluation rather than the whole answer. AutoMQ is a Kafka-compatible cloud-native streaming platform built on a Shared Storage architecture. It keeps Kafka-facing clients and protocol semantics while separating broker compute from durable object storage through its S3Stream storage layer.
That architecture can be relevant when CDC retention and replay reads compete with live connector traffic. It gives a team another storage boundary to test on GCP, while the application contract stays the same: verify partition keys, offsets, connector state, schema behavior, sink idempotency, and recovery. The AutoMQ architecture overview explains the storage layers, and the GKE deployment guide covers the Google Cloud deployment boundary. The GCP Kafka replacement guide and GKE deployment pattern provide adjacent migration and infrastructure context.
AutoMQ does not remove the need to test CDC semantics. A Kafka-compatible cluster can still receive a wrongly keyed event, a duplicate retry, or an unsafe sink write. Its role is to let the team compare a shared-storage Kafka architecture against the operational and replay evidence gathered from the existing GCP path. Choose from that evidence.
7FAQ
7.1Does Kafka preserve the order of all CDC records?
Kafka preserves order within a partition. Preserve per-entity order by using a stable key and testing the source connector's partitioning behavior. Cross-partition global order is a separate requirement and usually reduces parallelism.
7.2Can a replay create duplicate database changes?
Yes. A task can fail after the sink accepts a record but before the consumer commits its offset. Use an idempotency key or version-aware upsert, then verify the result with a bounded replay test.
7.3Is a managed GCP Kafka service enough for CDC correctness?
A managed broker can cover broker operations, but connector offsets, source positions, schema changes, retries, and sink writes remain part of your application contract. Test the full path.
7.4What should we test first?
Pick one high-value aggregate, one bounded source or Kafka range, and one sink reconciliation query. Prove ordering, duplicate handling, and final state before widening the workload.
The first question was whether a CDC connector being “up” proves that a GCP Kafka pipeline is correct. It does not. Start with the source position and the sink's write rule, then use offsets, retries, and replay to make the contract observable. If you want to evaluate a Kafka-compatible shared-storage path on GCP, request an AutoMQ evaluation with the same CDC workload and acceptance tests.
