Table of Contents
Table of Contents
“If it speaks the Kafka protocol, do we still need to test Kafka Connect?”
Yes. A reachable bootstrap server proves that a Worker can open a client connection. It does not prove that the Worker can discover the metadata it needs, join its Connect group, distribute Tasks, commit and recover Offsets, load the same plugin dependency tree, or keep making progress when the workload changes shape. Those are separate contracts, and a production Connector exercises all of them in one run.
The first green produce-and-consume check is useful because it catches endpoint, authentication, and basic metadata mistakes. It is also where many evaluations stop. A Connector can fail later while the Broker remains reachable, and the failure can look like a plugin or application issue when the real difference is group coordination, internal Topic state, or a runtime boundary.
The acceptance question is narrower: can the existing Connect workload perform its required actions on the target, recover from the failures the team cares about, and provide evidence that its Records and progress are correct? The answer needs five acceptance layers, each with a minimum test and a pass line.
1Protocol compatibility is the starting line
A compatibility claim describes an important boundary, but it does not replace a workload test. The Apache Kafka Connect documentation describes source and sink Connectors running in standalone or distributed Workers. Distributed Workers use Kafka for coordination and persistent Connector state, so Connect is a demanding Kafka client: it creates Topics, joins groups, fetches metadata, commits Offsets, and repeatedly exercises paths that an application uses during a failure.
A useful test plan separates five surfaces:
- Connection and metadata: Can the Worker authenticate, discover Topics and Partitions, and perform the Kafka requests its Connector needs?
- Task assignment: Can Workers join the same Connect group, receive assignments, rebalance, and return to a stable state?
- Offset behavior: Can source or sink progress be committed, recovered, and checked at a known boundary after a restart?
- Plugin loading: Can the target runtime load the Connector class, Converters, Transforms, and dependencies with the intended configuration?
- Workload shape: Does the Connector keep making useful progress during steady traffic, bursts, replay, backpressure, and an external-system slowdown?
A test can pass one layer and fail another. “The Connector started” is therefore a weak acceptance statement. The stronger statement names the layer, the test, the evidence, and the boundary that remains open.
AutoMQ is the platform under test here. AutoMQ’s official Apache Kafka compatibility documentation describes its use of the Apache Kafka computing layer with storage-layer changes and lists compatibility with Kafka clients, Connectors, and other Kafka ecosystem components for relevant versions. That is a starting point, while the acceptance result belongs to the team’s Connector, configuration, and recovery evidence.
For teams using AutoMQ BYOC managed Kafka Connect, the official documentation describes Connect Workers deployed in the user’s VPC and separate plugin and Connector management paths. Those details change who operates the Worker. They do not remove the five tests.
2Test one: connection and metadata
This layer asks whether the Worker can use the target as the Kafka cluster it expects. Network reachability and authentication are necessary, but the test should continue through Topic discovery, record production or fetching, and Offset lookup. A source may create or describe a Topic; a sink may discover assigned Topics and read committed Offsets before its first batch.
The minimum test is a short source-to-sink path with one representative Connector and a test Topic with more than one Partition. Produce Records with keys, headers, and a recognizable payload. Capture Worker logs, Topic metadata, Partition count, record keys, and the first and last Offset observed by the sink. Keep any Schema Registry or external system in the test.
Pass when the Worker authenticates, discovers the expected metadata, reaches the resources it owns, and preserves the expected key, value, headers, and per-Partition ordering without an application-code workaround. This proves a usable Kafka boundary. It does not prove Task rebalancing or Offset recovery.
3Test two: Task assignment and group coordination
Distributed Kafka Connect Workers are members of a Connect Worker group. The group coordinates which Worker owns a Connector and its Tasks. A Connector that runs on one Worker may therefore appear healthy while a two-Worker deployment exposes a rebalance, assignment, or session-timing problem. The target cluster is involved because the Worker group uses Kafka requests and internal state, even though Connect owns the scheduler.
Start two Workers with the same group.id, Worker configuration, and plugin set. Deploy one Connector with enough Task capacity to produce more than one Task when the Connector supports it. Confirm the initial assignment, stop one Worker, wait for the rebalance, start it again, and record which Tasks move, whether the remaining Task continues, and how the Connector reports the transition. Include a source or sink with a visible side effect so liveness cannot be confused with progress.
Pass when the rebalance reaches a stable assignment within the team’s recovery budget, active Tasks avoid an unexplained duplicate or gap, and the Connector returns to normal progress after the Worker rejoins. A transient REBALANCING status is different from a group that never reaches assignment.
Capture membership and assignment, Connector and Task status, records before and after the Worker change, and rebalance timestamps. A producer and Consumer can work while the Connect group fails to converge, so test the group as a first-class dependency.
4Test three: Offset commit and recovery
Offsets mark the boundary between “the Connector has read this” and “the Connector can resume from here.” The meaning depends on the Connector. A sink usually commits Kafka Offsets for records it has processed according to its delivery behavior. A source persists a position in the source system and then writes Records to Kafka. State which side of the boundary the test measures.
Produce a known sequence to a test Topic, let the Connector process through a recorded boundary, stop the Worker or Task, and inspect the committed Offset and downstream side effect. Restart the Connector and compare the post-restart sequence with the expected records. Repeat with a controlled failure where the external write may have completed while its response was lost.
Pass when the last committed Offset is visible, restart resumes from the documented boundary, replay is explainable and safe for the destination, no record is silently skipped, and operators can distinguish replay from a destination or application duplicate.
Do not turn this into an unqualified exactly-once promise. The Kafka consumer configuration reference explains client-side controls around groups, commits, and position, while the destination can impose its own idempotency and transaction boundary. An Offset ledger with Topic, Partition, pre-failure Offset, post-restart Offset, and destination record ID gives the team evidence beyond a green status screen.
5Test four: plugin loading and runtime dependencies
A Connector is a package, not a class name. The Worker must find the implementation, Converter classes, Single Message Transforms, and every dependency the plugin requires. A package can be present and still fail because a dependency is missing, plugin libraries conflict, or the runtime uses different Converter settings from production.
Start a clean Worker with the exact plugin artifact, version, Connector class, Converter, Transform, and configuration intended for the target. Register the Connector, run a representative Record through the complete serialization path, and restart the Worker from the same package. For a source, include the source position and first emitted Record. For a sink, include encoding, headers, and destination response.
Pass when the plugin loads without classpath repair, the Record crosses the intended Converter and Transform path with the expected result, and a restart reproduces the result from the same package and configuration. Record the negative case too: a missing class should fail visibly with an actionable error rather than leave a Connector reporting RUNNING without a Task.
AutoMQ’s plugin documentation describes managed and custom plugin paths for its managed Kafka Connect service, including validation of custom packages before creation. Use the plugin management documentation to identify the supported path, then test the exact plugin and version your pipeline owns. Class loading still does not prove that the plugin can reach a database, object store, or API with the target identity.
6Test five: workload shape and useful progress
A Connector can pass every quiet-path test and still fail the workload it was hired to run. The difference is often shape: a burst after an outage, a long replay, a slow sink, uneven Partition traffic, larger Records, a schema change, or a source that pauses and resumes. These conditions change data in flight, poll timing, Offset commits, and pressure on the external system.
Use three windows: steady traffic, a burst or replay, and a controlled downstream slowdown. Keep the record shape representative, including key distribution, headers, serialization format, and Partition skew. Record input and output rate, Consumer lag, Task status, errors, retries, Worker resource use, and time to return to normal progress. The goal is behavior evidence, not a universal throughput claim.
| Window | Pass condition |
|---|---|
| Steady traffic | The Connector meets the team’s freshness or delivery objective, and lag does not grow without bound. |
| Burst or replay | Tasks keep making progress, backlog is visible, and recovery does not depend on an undocumented manual step. |
| Downstream slowdown | Backpressure is observable, retries follow policy, and the destination is not overwhelmed by an uncontrolled retry storm. |
Write down what the test did not cover. Uniform Partition traffic does not cover skew. A short run does not prove long-term stability. A source-only test does not cover sink delivery. The missing cases should become explicit follow-up gates, with an owner and a reason they are outside the first acceptance window. That keeps a narrow test from being mistaken for a broad performance statement. For related operating checks, see connector lifecycle readiness and connector fleet observability.
7The surprises: group coordination and internal Topics
Two failure patterns sit between the five layers. First, a Connector’s Tasks consume Kafka Records as members of a consumer group, while Connect Workers coordinate their own membership and assignments. A healthy data Consumer group does not prove a healthy Connect Worker group. Capture the group ID, session timing, poll behavior, assignment protocol, and rebalance events alongside data-path evidence.
Second, distributed Connect commonly uses Topics for configuration, source Offsets, and Task or Connector status. Their names and settings are Worker configuration, not decorative labels:
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-statusThe exact names are a team choice, but the state is part of the Connector runtime. These Topics need the permissions, retention, Partition layout, and availability expected by the Worker. If they are absent, inaccessible, shared with the wrong Connect cluster, or recreated without a deliberate reset decision, the symptom can be a Task failure, unexpected replay, missing status, or a Worker that never reaches stable assignment.
The minimum internal-state test is to deploy with explicitly recorded internal Topic settings, verify that the expected Topics exist and are writable, restart one Worker, and inspect the same config, Offset, and status state before and after the restart. Pass when the Worker reads the intended state, keeps the expected assignment, and resumes from the documented Offset boundary. Starting with empty internal Topics is a fresh-deployment test, not Offset preservation.
AutoMQ’s Connector management documentation calls out Kafka credentials for the Topics and consumer groups a Connector uses, plus Worker Configuration and source Offset settings. Those are useful checks for a managed deployment. They also show why permissions and state belong in the acceptance evidence.
8Build an acceptance gate that can say go
The gate should prevent one green test from speaking for five untested surfaces. Give each layer an owner, an executable test, a pass condition, an evidence location, and a failure case from the production runbook.
| Gate | Required evidence | Decision rule |
|---|---|---|
| Connection and metadata | Authenticated metadata and record-path capture | Hold if the Connector needs an application workaround for the Kafka boundary. |
| Task assignment | Two-Worker rebalance and stable assignment record | Hold if assignment does not converge or progress stops after a Worker change. |
| Offsets | Pre-failure and post-restart Offset ledger | Hold if a gap is unexplained or replay safety is undefined. |
| Plugin loading | Clean package, dependency, Converter, and Transform test | Hold if classpath repair or an untracked artifact is required. |
| Workload shape | Steady, burst or replay, and slowdown evidence | Hold if lag, retries, or recovery behavior has no owner or objective. |
Record the Apache Kafka client and Connect versions, plugin artifacts, Worker properties, internal Topic names, security settings, source and sink versions, record shape, Partition layout, and test date. Keep the target and baseline configurations side by side so a failure can be attributed to a platform difference instead of a changed test input. Name the person or team that owns each permission, plugin artifact, and recovery decision. If a test uses a managed service path, record which settings are inherited and which remain under team control. The same record shape should be used for the baseline and target, including null values and headers when the Connector depends on them. That makes the acceptance packet reusable when the plugin or Worker version changes. If AutoMQ BYOC managed Kafka Connect is selected, record the Connect Cluster, Worker count, Worker Tier, plugin version, and external-system permissions described in the official Connector management documentation.
An acceptance packet should contain a test matrix with a named owner, an evidence bundle with logs and Offset ledgers, and a decision record listing passed layers, open boundaries, rollback conditions, and the workload that remains untested. That tells the platform team what can move, what must be configured, and what still needs a workload-specific decision.
Protocol compatibility answers whether Connect can speak to the platform. Acceptance testing answers whether your Connector can keep its promises after it starts speaking. Run the five tests against one representative source or sink, keep the pass lines visible, and treat an unknown recovery boundary as a hold. When you are ready to run that matrix against your own workload on a Kafka-compatible streaming platform, start an AutoMQ evaluation.
