Table of Contents
Table of Contents
Kafka Connect on GCP can be healthy enough to authenticate and still be unable to make progress. The worker polls records, the destination accepts some requests, a quota starts rejecting others, and the connector spends its time retrying work that is already late. Kafka keeps the records, but the application sees growing lag and an unclear recovery path.
Google Cloud pipelines expose this mismatch because a sink and Kafka have different capacity boundaries. A BigQuery, object-storage, or service-API sink may enforce quotas, batch limits, or concurrency while Kafka continues to expose partitions. The design question is how should the pipeline slow down, preserve required ordering, and make a failed batch safe to replay?
1The sink is part of the delivery contract
Kafka Connect workers do not make a destination transactional by themselves. A source connector can publish records to Kafka while a sink connector consumes them and reports progress according to its task semantics. When a destination write fails, the connector may retry, pause, quarantine a record, or fail the task. The choice depends on what the destination acknowledges and what the application considers a duplicate.
Treat the sink contract as a set of explicit answers:
| Question | Decision to record | Evidence to keep |
|---|---|---|
| When is a record accepted? | Destination response, durable batch commit, or application-level acknowledgment | Request ID, partition, offset, and destination receipt |
| What is retryable? | Network errors, throttling, temporary unavailability, or a bounded error class | Error code and retry policy |
| What is not retryable? | Invalid schema, authorization failure, malformed payload, or a poison record | Failed record reference and owner |
| Where can order change? | Within a partition, within a batch, or across tasks | Ordering rule and replay result |
| What happens after a task restart? | Resume from a committed position, replay a batch, or quarantine it | Connector log and offset decision |
A connector that says “task is running” answers only a process question. The delivery contract answers whether the destination has accepted the data and whether a restart can explain every record that follows. For a managed Kafka deployment, Google’s Kafka Connect overview is the reference point for the service boundary; the failure contract still belongs to your workload.
2Follow three clocks instead of one lag number
Sink incidents become easier to reason about when the team watches three clocks. Kafka has an arrival clock: the rate at which producers append records. The connector has a processing clock: the rate at which tasks poll, transform, and submit batches. The destination has an acceptance clock: the rate at which it commits writes under its quota and error policy. Backpressure appears when the arrival clock stays ahead of either downstream clock.
The same lag metric can mean different things depending on which clock moved. A producer surge may call for capacity or an intake limit. Throttling may make more tasks harmful, while schema or authorization failures will not yield to parallelism.
Capture these signals together rather than paging on lag alone:
- Consumer lag by topic, partition, connector, and task, with the oldest record age when available.
- Batch size, request rate, retry count, backoff time, and task restarts.
- Destination response classes, quota signals, and accepted-versus-rejected records.
- Dead-letter volume, replay queue depth, and the identity of the last acknowledged batch.
The combination shows whether the connector is starved, throttled, or blocked on a record that cannot succeed. If the evidence points to partition or application limits rather than the sink boundary, the diagnostic split in Kafka Consumer Lag at Scale is a useful next check.
3Classify the failure before changing retries
A retry policy is a control decision, not a universal safety switch. Start by classifying the failure at the boundary where it occurs. The following worksheet keeps transient faults separate from records that require a human or schema change.
| Failure class | Typical symptom | Control action | Replay requirement |
|---|---|---|---|
| Transient network or service fault | Requests time out or return a temporary availability error | Use bounded backoff and preserve the batch for another attempt | Prove that the destination write is idempotent or deduplicated |
| Quota or rate limit | Throttling responses increase as concurrency rises | Reduce request pressure, coordinate with the quota owner, and watch recovery | Re-send only batches whose acceptance is known |
| Poison record or schema mismatch | The same record fails while later records wait | Quarantine with enough context to repair or discard safely | Reprocess the record from a controlled offset or dead-letter path |
| Identity or configuration failure | Authentication, permission, endpoint, or routing errors persist | Stop scaling tasks and repair the boundary condition | Verify the first successful write before releasing backlog |
Change one control surface at a time. Raising parallelism, increasing batch size, and extending retries together makes the result impossible to interpret. Keep the error, connector configuration, and destination response together in the incident record.
4Backpressure is a control loop
Backpressure works when it closes the loop between destination feedback and Kafka consumption. The connector observes a response, classifies it, changes how much work it offers downstream, and measures whether the destination recovers. A static retry count is only one input to that loop.
For a sink on GCP, design the loop around four boundaries:
- Batch boundary. Keep the unit of retry small enough that one rejected record does not hide the outcome of an entire backlog. The right size is determined by destination semantics and payload shape, not by a generic connector default.
- Concurrency boundary. Treat tasks and in-flight requests as pressure on the destination. More tasks can improve throughput when the destination has capacity, but they can also multiply throttling and make replay order harder to explain.
- Time boundary. Use backoff and retry timeouts that leave room for an operator to classify the fault. An unbounded retry can preserve data while silently consuming the recovery window.
- Quarantine boundary. Give non-retryable records a controlled destination, such as a dead-letter topic or an equivalent review queue, with the source topic, partition, offset, key, and error context preserved.
Kafka Connect exposes consumer and connector settings that influence polling, batching, retries, and error handling. Review the selected worker and connector configuration against the Kafka Connect configuration reference, then test the exact runtime and destination combination. A setting that is valid for one connector implementation does not automatically define the behavior of another.
The objective is to keep the pipeline in a state where the destination can recover and the operator can explain which records were accepted.
5Build a lag budget from destination capacity
A lag budget connects a business recovery expectation to measurable connector behavior. Start with the oldest record age or another workload-specific freshness target. Then reserve time for destination retries, operator response, and replay of quarantined records. The remaining room is the connector’s operating budget.
Do not turn this into a universal threshold. A stream feeding a near-real-time dashboard may have a tighter freshness contract than a batch-oriented warehouse load, while both use Kafka Connect. Record the contract with the destination quota and the alert that should fire when the budget is being consumed.
Alert on budget consumption rather than a single static lag value:
- Alert when oldest record age approaches the workload’s freshness boundary.
- Alert when retry time or throttled request share consumes the reserved recovery room.
- Alert when the accepted batch watermark stops advancing while producer traffic continues.
- Alert when dead-letter or replay volume grows without a confirmed owner and runbook.
If the destination quota is the limiting resource, a scaling event should reduce pressure, not amplify it. If the connector is CPU- or network-bound while the destination accepts requests normally, capacity changes may help. The evidence should identify which side of the boundary moved before the team changes topology.
6Test replay before an outage chooses for you
Replay is where a sink design proves its assumptions. A successful retry does not tell you whether a task restart can duplicate a batch or whether a dead-letter record can be repaired.
Use a bounded test topic and a destination dataset that can be inspected. Record a producer-generated event ID, the source partition and offset, the connector task, and the destination request or receipt identifier. Then run failure cases that match the contract:
- Return a temporary destination error and verify that backoff limits pressure while the accepted watermark remains explainable.
- Apply a quota response and confirm that adding concurrency is not the recovery action.
- Stop a task after submission but before its progress is committed, then measure duplicate handling.
- Send a non-retryable record and verify quarantine, repair, and controlled reprocessing.
- Restart the worker and compare event IDs, ordering, accepted batches, and destination state.
The pass condition is an evidence bundle, not a green connector status. Show which source positions were attempted, which writes were accepted, which records were quarantined, and how replay changed the destination. If the sink cannot distinguish “unknown acceptance” from “rejected,” add an idempotency key or destination reconciliation step.
7Where a shared durable stream changes the pressure boundary
The worksheet applies to managed Kafka on Google Cloud, GKE, and self-managed compute. The connector still consumes Kafka records and the destination still enforces its contract. Storage architecture changes how broker storage pressure and connector replay compete for local resources.
AutoMQ is a Kafka-compatible cloud-native streaming platform built around a Shared Storage architecture. Its S3Stream layer keeps durable stream data in object storage, with a WAL (Write-Ahead Log) layer and broker caches handling write and read paths. The AutoMQ architecture overview describes those storage and broker boundaries. That separation can make the connector’s replay workload less coupled to broker-local data placement, while the connector continues to use Kafka protocol and consumer semantics.
This is an architectural hypothesis, not a promise that a sink ignores quotas. Run the same worksheet against the selected AutoMQ deployment and measure lag, retry pressure, broker resource use, destination acceptance, and replay outcomes. Record the GCP network path, IAM identities, object-storage location, WAL type, connector runtime, and sink configuration so the result is reproducible.
8FAQ
8.1Is Kafka Connect on GCP responsible for destination backpressure?
The connector participates in the control loop, but the destination defines the acceptance boundary. A sound design uses destination responses to adjust batch pressure and task behavior, while Kafka retains records for the configured retention and replay policy.
8.2Should I add more connector tasks when sink lag grows?
Only after the evidence shows that the destination can accept more concurrency. If throttling or quota responses are increasing, more tasks can multiply pressure. Classify the failure and check the lag budget first.
8.3How should a BigQuery Kafka sink handle a failed batch?
Define what the destination acknowledges, preserve source partition and offset metadata, and test duplicate handling after a task restart. Separate transient service or quota errors from schema, identity, and malformed-record failures; the recovery path is different for each.
8.4What should a GCP Kafka sink replay test prove?
It should connect source positions to destination outcomes. Capture event IDs, offsets, request or receipt identifiers, quarantined records, and the final destination state. A connector that restarts successfully without this evidence has proven process recovery, not data delivery.
A sink incident starts with a mismatch between Kafka’s arrival clock and the destination’s acceptance clock. The way out is to make that mismatch visible, bound retries, and rehearse replay. If you want to test the same connector contract on a Kafka-compatible shared-storage platform, start with AutoMQ Open Source.
