Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Connect Exactly-Once

Kafka Connect supports exactly-once semantics (EOS) for both source and sink connectors.


Available in Kafka 3.3+. Source connectors use transactions to ensure exactly-once delivery.

# Enable exactly-once support
exactly.once.source.support=enabled
# Transaction boundary mode
transaction.boundary=poll
Transaction BoundaryBehavior
pollTransaction per poll() call (default)
connectorConnector defines boundaries
intervalTransaction every N milliseconds
Source TaskKafkaExternalSource TaskSource TaskKafka(Transactional)Kafka(Transactional)ExternalSourceExternalSourcepoll()recordsbeginTransaction()produce(records)sendOffsetsToTransaction()commitTransaction()Records and offsetscommitted atomically

No additional connector config needed; EOS uses worker settings:

{
"name": "eos-source",
"config": {
"connector.class": "...",
"tasks.max": "1"
}
}

Sink connectors achieve exactly-once through idempotent writes, not Kafka transactions.

StrategyImplementation
UpsertPrimary key ensures idempotent updates
DeduplicationTrack offsets in sink system
External transactionsCommit offset with sink transaction

Use natural keys for idempotent writes:

{
"name": "cassandra-sink",
"config": {
"connector.class": "com.datastax.oss.kafka.sink.CassandraSinkConnector",
"topics": "events",
"topic.events.keyspace.table.mapping": "analytics.events"
}
}

Cassandra’s upsert semantics ensure duplicates overwrite with same data.

Store Kafka offsets alongside data:

-- PostgreSQL example
CREATE TABLE events (
kafka_topic VARCHAR(255),
kafka_partition INT,
kafka_offset BIGINT,
event_data JSONB,
PRIMARY KEY (kafka_topic, kafka_partition, kafka_offset)
);

AspectSource EOSSink EOS
MechanismKafka transactionsIdempotent writes
Kafka version3.3+Any
ConfigurationWorker-levelConnector/sink-level
OverheadTransaction coordinationSink-dependent

LimitationDescription
Source connector supportConnector must be EOS-compatible
Sink external systemsMust support idempotent writes
PerformanceTransaction overhead on source side
Zombie fencingRequires proper task assignment

Check connector status for EOS mode:

Terminal window
curl http://connect:8083/connectors/my-source/status | jq '.tasks[].trace'