Blog

Diskless Kafka Kafka Connect: A Production Framework

Table of Contents

Table of Contents

Kafka Connect can be healthy while the pipeline it runs is not. Workers may be online, connector tasks may be assigned, and the target system may still be missing records because a sink task is replaying, a worker cannot read its internal topics, or a cold fetch from object storage is taking longer than the task’s poll loop expects. The operational question is therefore bigger than “does Kafka Connect support diskless Kafka?” It is whether the team can prove offset continuity, restart behavior, and target-system recovery across the storage paths that the connector actually uses.

That proof starts with a boundary. Kafka Connect remains a distributed runtime made of workers, connectors, and tasks. A diskless Kafka design changes where Kafka retains durable stream data; it does not remove the Connect worker’s configuration, status, offset, group, plugin, network, and credential dependencies. Treat those as separate contracts, then test the byte path that joins them.

1What Kafka Connect means in a diskless Kafka design

Kafka Connect has two traffic directions. A source connector reads from an external system and produces records to Kafka. A sink connector consumes Kafka records and writes them to an external system. The connector task is the unit that does the work, but the worker cluster decides where tasks run, stores connector state, and coordinates changes. The Apache Kafka Connect documentation describes the standalone and distributed modes, the worker configuration, and the internal topics used by distributed deployments.

Those internal topics are part of the production data path even when the connector’s business data is elsewhere. Distributed workers use Kafka topics for connector configuration, task status, and offsets; the exact offset behavior also depends on whether the connector is a source or a sink and on the connector implementation. A worker restart is safe only when the replacement worker can reach the same Kafka cluster, authenticate, load a compatible plugin, read the relevant state, and rejoin the worker group. Moving Kafka’s durable data to object storage does not move these responsibilities into an abstract control plane.

This is the first useful distinction for a diskless Kafka evaluation:

BoundaryWhat must remain trueEvidence to retain
Connect workerWorkers can start with the same plugin, configuration, credentials, and network reachabilityStartup logs, plugin version, resolved configuration, and worker membership
Connect stateConfiguration, task status, and source or sink progress remain readable and writableInternal-topic health, partition leadership, committed positions, and state transitions
Kafka dataThe connector can fetch or append the records it owns through the selected storage pathProduce or fetch latency, bytes, cache state, storage reads, and errors
External systemThe source or target accepts the connector’s requests and preserves its own idempotency or transaction contractRequest results, retry behavior, duplicate handling, and application-level checkpoints

The table separates a worker failure from a storage failure. That separation matters during an incident. A sink task that is waiting for a cold Kafka range can look like a target-system outage, while a target-system throttle can make Kafka lag look like an object-storage problem. The investigation has to follow the task’s state and the record’s bytes at the same time.

2Follow the connector’s data path

In a broker-local Kafka deployment, a sink task commonly reads recent records from the broker’s local log and page cache. A replay that moves behind the local working set may hit a remote tier when the deployment uses Kafka Tiered Storage. KIP-405 describes that local-tier and remote-tier model. A diskless Kafka design puts shared storage at the durable boundary for the stream, so the broker can serve hot data from cache while a catch-up read reaches object storage through the storage layer selected by the deployment.

The connector does not see that architecture as a separate application API. It still sends Kafka protocol requests and receives records, errors, and offsets through the client boundary. The operational difference is inside the path: a broker replacement, a cache miss, an object-store request, or a network boundary can now change how the same fetch is served. The test plan should make that hidden path visible.

For a sink connector, trace one record batch through these points:

  1. The task polls Kafka and receives a response with the expected partitions and offsets.
  2. The broker serves the range from its hot cache, a write buffer, or a catch-up path that reads shared storage.
  3. The task transforms the records and sends them to the external target.
  4. The target acknowledges the write according to its own contract.
  5. The task commits progress through the Kafka Connect and consumer-group mechanisms that apply to that connector.

For a source connector, reverse the direction. The task reads the source position, produces to Kafka, waits for the producer contract that the connector configures, and persists the source position so a restart can resume at the intended point. A diskless Kafka test must cover both directions when the deployment runs both source and sink workloads; a sink-only test cannot prove source offset recovery.

The storage mechanism becomes a practical concern when a task falls behind. A tailing read may stay close to cache, while replaying a wide offset range may produce more object-storage requests, cache evictions, and network traffic. The connector’s poll and batch settings then interact with the storage path: small fetches can increase request overhead, while large batches can increase memory pressure or make a target timeout more likely. Record the settings and the observed path instead of assuming that a connector’s default behavior is representative of production.

AutoMQ describes this shared-storage path as Kafka protocol handling on brokers, S3Stream for streaming storage, a WAL layer for durable writes and recovery, and a data cache for hot and catch-up reads. The architecture is useful here as a measurement model: a Connect task still sees Kafka semantics, while the operator can attribute latency to broker scheduling, cache behavior, WAL, object storage, or the network between them. The AutoMQ architecture overview provides the component boundaries to map in a pilot.

3Failure, cost, and compatibility checks

A connector restart is not one test. It is a chain of recovery decisions, and each decision needs a pass condition. Start with the failure that the team already rehearses, then add the storage and target failures that a local-disk mental model tends to hide.

DrillWhat to observePass condition to define before the test
Restart one Connect workerTask reassignment, plugin loading, internal-topic reads, and progress after assignmentThe task returns to the intended worker group and resumes according to the connector’s offset contract
Replace or restart a Kafka brokerFetch or produce path, leader readiness, cache warm-up, and task poll behaviorThe connector either continues or pauses within the declared operating policy, with no unexplained offset jump
Replay from an older Kafka offsetCache misses, object-store reads, request retries, and task processing timeThe task reaches the replay target with a measured recovery path and a known backpressure response
Add object-storage latency or errorsFetch latency, retries, lag, worker health, and alert ownershipThe team can distinguish storage degradation from worker or target failure and take the documented action
Throttle or stop the external targetSink retries, task state, consumer progress, duplicate handling, and recovery after the target returnsThe connector’s delivery and idempotency contract remains explicit; recovery does not depend on an operator guessing the last safe offset
Rebalance workers or tasksAssignment changes, state reads, task startup, and connector-level orderingThe replacement assignment reaches a stable state and the resulting data path is observable

The pass condition should use the connector’s own contract. “No data loss” is not a complete test statement unless the team defines which records were acknowledged, which offsets were committed, and how the target handles retries. Exactly-once behavior, ordering, and duplicate suppression are connector- and target-specific; object storage does not make any of them automatic. Keep the test result tied to the connector, Kafka client, target, and configuration that produced it.

Cost belongs in the same worksheet because a replay can change the bill without changing retained bytes. Separate the lines for Connect compute, Kafka broker compute, object-storage capacity, object-store requests, WAL storage where applicable, network transfer, and observability. Then mark which lines grow during tailing, replay, connector restart, or target backpressure. A statement such as “Kafka on S3 changes the storage line” hides the request and network pattern that the connector creates.

Compatibility needs the same evidence discipline. Inventory every plugin and its dependency tree, the Kafka client version it carries, authentication and authorization requirements, schema or converter behavior, transaction settings, compaction assumptions, and the target’s retry and idempotency rules. The KIP-1150 Diskless Topics proposal is useful context for the diskless topic concept, but it is a proposal reference rather than proof that a chosen Kafka release or connector deployment implements the behavior you need. Validate the release, plugin, and storage configuration together.

The most useful evidence is partition-specific. Keep the task assignment, consumer position, latest partition offset, fetch response size, fetch latency, processing time, commit result, cache state, storage-read latency, and target response in one timeline. Group averages can hide a single hot partition or a task that is repeatedly replaying the same range. If the timeline cannot explain why a task stopped making progress, the production gate is not ready.

4How AutoMQ changes the operating model

Once the worksheet shows that local-disk ownership is driving recovery or scaling work, a Kafka-compatible Shared Storage architecture becomes a concrete option. AutoMQ preserves the Kafka protocol boundary used by Connect while placing durable stream data in S3-compatible object storage. That lets a team evaluate broker replacement, cache behavior, WAL choice, and object-store access with the same connector tests instead of changing the connector application first.

The architecture changes the operator’s question. A broker replacement no longer means that the replacement must first own a complete local copy of every retained partition before a task can fetch it. The replacement still needs the right metadata, network path, permissions, cache behavior, and WAL configuration. Shared storage changes where durable bytes are found; it does not remove the need to measure the path from a Connect poll to an external acknowledgement.

WAL choice should be explicit in the pilot. AutoMQ documentation describes S3 WAL, Regional EBS WAL, and NFS WAL as deployment options with different latency and failure-domain properties. The WAL storage documentation should be read alongside the target cloud and network design. Record where the WAL and cache sit, which failures they are expected to absorb, and which alert tells the Connect operator that the dependency is degraded.

The Connect cluster remains an independent operational surface. Decide whether workers run beside the brokers, in a separate node pool, or as a managed service. In each case, test worker-to-broker reachability, plugin distribution, credentials, internal-topic placement, and target-system access. A diskless Kafka broker can be stateless with respect to retained stream data while a Connect deployment still depends on stable worker configuration and state topics.

That distinction is useful during migration. Move one connector or one production-shaped flow, preserve its task and offset evidence, and compare the same tailing, replay, broker-replacement, and target-outage drills. If the candidate architecture removes local-log movement but adds object-store request pressure, the worksheet should show the resulting owner and the action it requires. The goal is a measurable operating model, not a storage label.

5Decision checklist and FAQ

Use these gates before expanding a diskless Kafka and Kafka Connect pilot:

  • State contract: The team has named the configuration, status, offset, and group state that each connector uses, along with its owner and backup or recovery procedure.
  • Path evidence: Every test records task assignment, partition offsets, fetch or produce behavior, cache state, storage reads, network path, target responses, and commit results in one timeline.
  • Failure policy: Worker loss, broker loss, object-storage degradation, network loss, and target throttling have explicit pause, retry, backpressure, and escalation actions.
  • Compatibility proof: The actual connector plugin, dependency set, converter, authentication method, Kafka release, and target system pass the same restart and replay scenarios planned for production.
  • Cost model: Compute, storage, object-store requests, network transfer, WAL, and observability are separate lines, with tailing and replay patterns measured independently.
  • Rollback gate: The migration has a stop condition, an offset strategy, a target reconciliation procedure, and an operator who has rehearsed the rollback.

5.1Does diskless Kafka change the Kafka Connect API?

The connector still uses the Kafka protocol and Connect worker model. The storage architecture changes how brokers serve durable stream data, so fetch, replay, recovery, and network behavior must be tested. API compatibility alone does not prove that a connector’s restart or target delivery contract is intact.

5.2Where do Kafka Connect offsets live in a diskless design?

They remain part of the Connect and Kafka state model selected by the deployment. Distributed workers use Kafka-backed internal topics for connector configuration, status, and offsets, while sink progress also involves consumer-group semantics and connector-specific behavior. Identify the exact topics, partitions, and commits for each connector and verify them during a worker and broker restart.

5.3Does object storage remove the need to size Connect workers?

No. Worker CPU, heap, task parallelism, plugin dependencies, network bandwidth, and target-side limits still bound the pipeline. Shared storage may change broker scaling and replay behavior, but a worker that cannot deserialize, transform, or deliver records will still fall behind.

5.4Is diskless Kafka the same as Kafka Tiered Storage?

They can both use object storage, but they place durable data differently. Tiered Storage commonly keeps an active local tier and moves older segments to a remote tier. Diskless Kafka makes shared storage the durable boundary for the topic. Confirm the release and implementation you are evaluating before assuming that a connector replay follows the same path in both designs.

5.5What should a production pilot prove first?

Start with one source or sink flow that has a clear freshness and recovery objective. Prove worker restart, broker replacement, cold replay, object-storage degradation, target outage, task rebalance, offset continuity, and rollback with the real plugin and target. Expand only when the evidence identifies the owner and action for every failed dependency.

Kafka Connect does not make a diskless Kafka design safe by itself, and diskless storage does not make Connect state disappear. The production decision becomes clear when one connector timeline can show which worker owned the task, which partition range it requested, which layer served the bytes, and when the external system acknowledged them. Run that timeline against your current cluster and a Kafka-compatible Shared Storage candidate, then start an AutoMQ evaluation with the failure gates and rollback conditions in hand.

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.