Skip to content

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

Kafka Connect

Kafka Connect is a framework for streaming data between Apache Kafka and external systems using pre-built or custom connectors.


Kafka Connect eliminates the need to write custom integration code for common data sources and sinks. The framework handles:

  • Parallelization and scaling
  • Offset management and exactly-once delivery
  • Schema integration with Schema Registry
  • Fault tolerance and automatic recovery
  • Standardized monitoring and operations
Event SourcesKafka Connect ClusterKafka ClusterData SinksREST APIsMQTT DevicesLog FilesWorker 1Worker 2Worker 3TopicsS3CassandraElasticsearchHTTP SourceMQTT SourceFileStreamS3 SinkCassandra SinkES Sink

ComponentDescription
WorkerJVM process that executes connectors and tasks
ConnectorPlugin that defines how to connect to external system
TaskUnit of work; connectors are divided into tasks for parallelism
ConverterSerializes/deserializes data between Connect and Kafka
TransformModifies records in-flight (Single Message Transforms)
Connect WorkerHerderConnectors & TasksConvertersREST API(:8083)Offset StorageConnectorManagementTaskManagementSource Task 1Source Task 2Sink Task 1KeyConverterValueConvertermanagecoordinateserialize/deserializetrack position
Source Connector FlowSink Connector FlowExternalSystemSourceTaskConverterTransform(SMT)KafkaTopicKafkaTopicTransform(SMT)ConverterSinkTaskExternalSystempollSourceRecordserializedproduceconsumetransformedSinkRecordwrite

Single worker process—suitable for development and simple use cases.

Terminal window
# Start standalone worker
connect-standalone.sh \
config/connect-standalone.properties \
config/file-source.properties \
config/file-sink.properties

Standalone properties:

bootstrap.servers=localhost:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
# Offset storage (local file)
offset.storage.file.filename=/tmp/connect.offsets
# REST API
rest.port=8083
CharacteristicStandalone Mode
WorkersSingle process
Offset storageLocal file
Fault toleranceNone
ScalingNot supported
Use caseDevelopment, testing

Multiple workers forming a cluster—required for production.

Terminal window
# Start distributed worker (on each node)
connect-distributed.sh config/connect-distributed.properties

Distributed properties:

bootstrap.servers=kafka1:9092,kafka2:9092,kafka3:9092
# Group coordination
group.id=connect-cluster
# Offset storage (Kafka topics)
offset.storage.topic=connect-offsets
offset.storage.replication.factor=3
offset.storage.partitions=25
# Config storage
config.storage.topic=connect-configs
config.storage.replication.factor=3
# Status storage
status.storage.topic=connect-status
status.storage.replication.factor=3
status.storage.partitions=5
# Converters
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://schema-registry:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081
# REST API
rest.advertised.host.name=connect-worker-1
rest.port=8083
CharacteristicDistributed Mode
WorkersMultiple processes (cluster)
Offset storageKafka topic (connect-offsets)
Fault toleranceAutomatic task redistribution
ScalingAdd workers to scale
Use caseProduction
TopicPurposeRecommended Config
connect-offsetsSource connector offsetsRF=3, partitions=25
connect-configsConnector configurationsRF=3, partitions=1, compacted
connect-statusConnector/task statusRF=3, partitions=5, compacted

Single worker process for development and simple integrations.

Single HostConnect Worker(Standalone)Kafka ClusterOffset File(/tmp/connect.offsets)Connector ATask 1TopicsExternalSystem- No fault tolerance- No scaling- Dev/test onlypersist offsetsproduce/consumeread/write

Multi-worker cluster for production deployments.

Connect ClusterWorker 1Worker 2Worker 3Kafka ClusterConnector ATask 1Connector BTask 1Connector ATask 2Connector CTask 1Connector ATask 3Connector BTask 2Data Topicsconnect-offsetsconnect-configsconnect-statusLoad Balancer(:8083)Internal topics provide:- Distributed offset storage- Configuration sync- Status coordinationREST APIREST APIREST API

Connect workers on Kafka broker nodes—suitable for smaller clusters.

Node 1Node 2Node 3Kafka BrokerConnect WorkerKafka BrokerConnect WorkerKafka BrokerConnect WorkerPros:- Lower latency (localhost)- Fewer nodes to manage Cons:- Resource contention- Failure affects both services- Harder to scale independentlylocalhostlocalhostlocalhostreplicationreplicationreplication

Separate Connect cluster—recommended for production.

Connect ClusterKafka ClusterExternal SystemsConnectWorker 1ConnectWorker 2ConnectWorker 3Broker 1Broker 2Broker 3DatabasesObject StorageAPIsPros:- Independent scaling- Isolated failures- Dedicated resources- Easier capacity planning

Connect worker per application—useful for application-specific integrations.

Application Pod 1Application Pod 2Application Pod 3Kafka ClusterApplicationConnectWorkerApplicationConnectWorkerApplicationConnectWorkerTopicsUse cases:- Application-owned connectors- Isolation requirements- Local file ingestion Note: Each can be standaloneor join distributed clusterlocal configlocal configlocal config

Connect cluster on Kubernetes with horizontal scaling.

Kubernetes ClusterConnect DeploymentServiceconnect-svc:8083IngressConfigMapconnect-configSecretconnect-secretsPod 1Connect WorkerPod 2Connect WorkerPod 3Connect WorkerKafka(internal or external)HPA can scale workersbased on CPU/memoryor custom metricsmountmount

Kubernetes manifest example:

apiVersion: apps/v1
kind: Deployment
metadata:
name: kafka-connect
spec:
replicas: 3
selector:
matchLabels:
app: kafka-connect
template:
metadata:
labels:
app: kafka-connect
spec:
containers:
- name: connect
image: confluentinc/cp-kafka-connect:7.5.0
ports:
- containerPort: 8083
env:
- name: CONNECT_BOOTSTRAP_SERVERS
value: "kafka:9092"
- name: CONNECT_GROUP_ID
value: "connect-cluster"
- name: CONNECT_REST_ADVERTISED_HOST_NAME
valueFrom:
fieldRef:
fieldPath: status.podIP
resources:
requests:
memory: "2Gi"
cpu: "1"
limits:
memory: "4Gi"
cpu: "2"
---
apiVersion: v1
kind: Service
metadata:
name: kafka-connect
spec:
ports:
- port: 8083
selector:
app: kafka-connect

Separate Connect clusters per region for geo-distributed workloads.

Region A (US-East)Connect Cluster ARegion B (EU-West)Connect Cluster BKafka Cluster ALocal SystemsWorkerWorkerKafka Cluster BLocal SystemsWorkerWorkerEach region has independent:- Connect cluster- Kafka cluster- Local integrations Cross-region via MM2MirrorMaker 2(cross-region)
PatternScalingFault ToleranceResource IsolationUse Case
StandaloneNoneNoneN/ADevelopment
Co-locatedWith brokersSharedPoorSmall clusters
Dedicated clusterIndependentIndependentGoodProduction
SidecarPer applicationPer applicationExcellentApp-specific
KubernetesHPA/VPAPod replacementGoodCloud-native
Multi-regionPer regionRegionalExcellentGlobal deployments
WorkloadWorkersMemory per WorkerCPU per Worker
Light (< 10 connectors)2-32-4 GB1-2 cores
Medium (10-50 connectors)3-54-8 GB2-4 cores
Heavy (50+ connectors)5-10+8-16 GB4-8 cores

Worker Sizing

Worker memory depends on:

  • Number of tasks per worker
  • Message size and throughput
  • Converter complexity (Avro/Protobuf vs JSON)
  • Transform chain depth

Terminal window
# Create a Cassandra Sink connector
curl -X POST http://connect:8083/connectors \
-H "Content-Type: application/json" \
-d '{
"name": "cassandra-sink",
"config": {
"connector.class": "com.datastax.oss.kafka.sink.CassandraSinkConnector",
"tasks.max": "1",
"topics": "events",
"contactPoints": "cassandra1,cassandra2,cassandra3",
"loadBalancing.localDc": "datacenter1",
"port": "9042",
"topic.events.keyspace.table.mapping": "events.events_by_time",
"topic.events.keyspace.table.consistencyLevel": "LOCAL_QUORUM"
}
}'
PropertyDescription
nameUnique connector name
connector.classFully qualified connector class
tasks.maxMaximum number of tasks
key.converterKey converter (overrides worker default)
value.converterValue converter (overrides worker default)
transformsComma-separated list of transforms
errors.toleranceError handling: none, all
errors.deadletterqueue.topic.nameDead letter queue topic
UNASSIGNEDRUNNINGPAUSEDFAILEDcreatetask assignedrebalancepauseresumeerrorrestartdeletedeletedelete

EndpointMethodDescription
/connectorsGETList all connectors
/connectorsPOSTCreate connector
/connectors/{name}GETGet connector info
/connectors/{name}DELETEDelete connector
/connectors/{name}/configGETGet connector config
/connectors/{name}/configPUTUpdate connector config
/connectors/{name}/statusGETGet connector status
/connectors/{name}/restartPOSTRestart connector
/connectors/{name}/pausePUTPause connector
/connectors/{name}/resumePUTResume connector
/connectors/{name}/stopPUTStop connector (deallocate resources)
EndpointMethodDescription
/connectors/{name}/tasksGETList tasks
/connectors/{name}/tasks/{id}/statusGETGet task status
/connectors/{name}/tasks/{id}/restartPOSTRestart task

Manage connector offsets (connector must be stopped):

EndpointMethodDescription
/connectors/{name}/offsetsGETGet current offsets
/connectors/{name}/offsetsDELETEReset offsets
/connectors/{name}/offsetsPATCHAlter offsets

Alter source connector offsets:

Terminal window
# Stop connector first
curl -X PUT http://connect:8083/connectors/my-connector/stop
# Alter offsets
curl -X PATCH http://connect:8083/connectors/my-connector/offsets \
-H "Content-Type: application/json" \
-d '{
"offsets": [
{
"partition": {"filename": "test.txt"},
"offset": {"position": 30}
}
]
}'
# Resume connector
curl -X PUT http://connect:8083/connectors/my-connector/resume

Reset sink connector offsets:

Terminal window
# Stop and reset to re-consume from beginning
curl -X PUT http://connect:8083/connectors/my-sink/stop
curl -X DELETE http://connect:8083/connectors/my-sink/offsets
curl -X PUT http://connect:8083/connectors/my-sink/resume
EndpointMethodDescription
/GETCluster info
/connector-pluginsGETList installed plugins
/connector-plugins/{plugin}/config/validatePUTValidate config

Dynamically adjust log levels:

EndpointMethodDescription
/admin/loggersGETList loggers with explicit levels
/admin/loggers/{name}GETGet logger level
/admin/loggers/{name}PUTSet logger level
Terminal window
# Enable debug logging for connector
curl -X PUT http://connect:8083/admin/loggers/org.apache.kafka.connect \
-H "Content-Type: application/json" \
-d '{"level": "DEBUG"}'

Converters serialize and deserialize data between Connect's internal format and Kafka.

ConverterFormatSchema Support
JsonConverterJSONOptional (schemas.enable)
AvroConverterAvro binaryYes (Schema Registry)
ProtobufConverterProtobuf binaryYes (Schema Registry)
JsonSchemaConverterJSON with schemaYes (Schema Registry)
StringConverterPlain stringNo
ByteArrayConverterRaw bytesNo
# JSON without schemas (simple)
key.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
# Avro with Schema Registry (production)
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://schema-registry:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081
Use CaseRecommended Converter
Development/debuggingJsonConverter (schemas.enable=false)
Production with schema evolutionAvroConverter or ProtobufConverter
Existing JSON consumersJsonConverter or JsonSchemaConverter
Maximum compatibilityStringConverter (manual serialization)

Converters Guide


SMTs modify records as they flow through Connect—useful for simple transformations without custom code.

TransformDescription
InsertFieldAdd field with static or metadata value
ReplaceFieldRename, include, or exclude fields
MaskFieldReplace field value with valid null
ValueToKeyCopy fields from value to key
ExtractFieldExtract single field from struct
SetSchemaMetadataSet schema name and version
TimestampRouterRoute to topic based on timestamp
RegexRouterRoute to topic based on regex
FlattenFlatten nested structures
CastCast field to different type
HeaderFromCopy field to header
InsertHeaderAdd static header
DropHeadersRemove headers
FilterDrop records matching predicate
{
"name": "my-connector",
"config": {
"connector.class": "...",
"transforms": "addTimestamp,route",
"transforms.addTimestamp.type": "org.apache.kafka.connect.transforms.InsertField$Value",
"transforms.addTimestamp.timestamp.field": "processed_at",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "(.*)_raw",
"transforms.route.replacement": "$1_processed"
}
}
SMT ChainOriginalRecordTransform 1(InsertField)Transform 2(ReplaceField)Transform 3(RegexRouter)FinalRecordTransforms execute in orderEach receives output of previousrecordmodifiedmodifiedmodified

Apply transforms conditionally based on message properties:

Built-in Predicates:

PredicateDescription
TopicNameMatchesMatch records where topic name matches regex
HasHeaderKeyMatch records with specific header key
RecordIsTombstoneMatch tombstone records (null value)

Predicate Configuration:

{
"name": "my-connector",
"config": {
"connector.class": "...",
"transforms": "FilterFoo,ExtractBar",
"transforms.FilterFoo.type": "org.apache.kafka.connect.transforms.Filter",
"transforms.FilterFoo.predicate": "IsFoo",
"transforms.ExtractBar.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.ExtractBar.field": "other_field",
"transforms.ExtractBar.predicate": "IsBar",
"transforms.ExtractBar.negate": "true",
"predicates": "IsFoo,IsBar",
"predicates.IsFoo.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.IsFoo.pattern": "foo",
"predicates.IsBar.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
"predicates.IsBar.pattern": "bar"
}
}
PropertyDescription
predicateAssociate predicate alias with transform
negateInvert predicate match (apply when NOT matched)

Transforms Guide


{
"name": "my-connector",
"config": {
"connector.class": "...",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "my-connector-dlq",
"errors.deadletterqueue.topic.replication.factor": 3,
"errors.deadletterqueue.context.headers.enable": true,
"errors.log.enable": true,
"errors.log.include.messages": true
}
}
ConfigurationDescription
errors.tolerance=noneFail on first error (default)
errors.tolerance=allLog errors and continue
errors.deadletterqueue.topic.nameTopic for failed records
errors.deadletterqueue.context.headers.enableInclude error context in headers
errors.log.enableLog errors to Connect log
errors.log.include.messagesInclude record content in logs
Sink ConnectorTaskSource TopicDead Letter QueueExternal SystemFailed records include:- Original key/value- Error message header- Exception stack trace- Source topic/partition/offsetconsumewrite (success)write (failure)

Error Handling Guide


Kafka Connect supports exactly-once semantics for source connectors (Kafka 3.3+).

# Worker configuration
exactly.once.source.support=enabled
transaction.boundary=poll # or connector, interval
# Connector configuration (automatically uses transactions)
Transaction BoundaryBehavior
pollTransaction per poll() call
connectorConnector defines boundaries
intervalTransaction every N milliseconds

Sink connectors achieve exactly-once through idempotent writes to external systems:

StrategyImplementation
UpsertUse primary key for idempotent updates
DeduplicationTrack processed offsets in sink
TransactionsCommit offset with sink transaction

Exactly-Once Guide


MetricDescriptionAlert Threshold
connector-countNumber of connectorsExpected count
task-countNumber of running tasksExpected count
connector-startup-failure-totalConnector startup failures> 0
task-startup-failure-totalTask startup failures> 0
source-record-poll-totalRecords polled by sourceDepends on workload
sink-record-send-totalRecords sent by sinkDepends on workload
offset-commit-failure-totalOffset commit failures> 0
deadletterqueue-produce-totalRecords sent to DLQ> 0 (investigate)
kafka.connect:type=connector-metrics,connector={connector}
kafka.connect:type=connector-task-metrics,connector={connector},task={task}
kafka.connect:type=source-task-metrics,connector={connector},task={task}
kafka.connect:type=sink-task-metrics,connector={connector},task={task}

Operations Guide


Add workers to the Connect cluster to distribute load:

Before ScalingWorker 1After ScalingWorker 1Worker 2Task 1Task 2Task 3Task 4Task 1Task 2Task 3Task 4Tasks automaticallyredistributed onworker join/leave

Increase tasks.max for connectors that support parallelism:

{
"name": "jdbc-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "10",
"table.whitelist": "orders,customers,products"
}
}
Connector TypeParallelism Model
HTTP SourceOne task per endpoint (typically)
File SourceOne task per file or directory
S3 SinkTasks share topic partitions
Cassandra SinkTasks share topic partitions

Building custom connectors to integrate Kafka with proprietary or unsupported systems.

ComponentInterfacePurpose
SourceConnectorSourceConnectorConfiguration and task distribution for imports
SourceTaskSourceTaskRead data from external system
SinkConnectorSinkConnectorConfiguration and task distribution for exports
SinkTaskSinkTaskWrite data to external system
public class MySourceConnector extends SourceConnector {
private Map<String, String> configProps;
@Override
public void start(Map<String, String> props) {
this.configProps = props;
}
@Override
public Class<? extends Task> taskClass() {
return MySourceTask.class;
}
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
// Distribute work across tasks
List<Map<String, String>> configs = new ArrayList<>();
for (int i = 0; i < maxTasks; i++) {
Map<String, String> taskConfig = new HashMap<>(configProps);
taskConfig.put("task.id", String.valueOf(i));
configs.add(taskConfig);
}
return configs;
}
@Override
public void stop() {
// Clean up resources
}
@Override
public ConfigDef config() {
return new ConfigDef()
.define("connection.url", Type.STRING, Importance.HIGH, "Connection URL");
}
@Override
public String version() {
return "1.0.0";
}
}
public class MySourceTask extends SourceTask {
private String connectionUrl;
@Override
public void start(Map<String, String> props) {
connectionUrl = props.get("connection.url");
// Initialize connection
}
@Override
public List<SourceRecord> poll() throws InterruptedException {
List<SourceRecord> records = new ArrayList<>();
// Read from external system
List<DataItem> items = fetchData();
for (DataItem item : items) {
Map<String, ?> sourcePartition = Collections.singletonMap("source", connectionUrl);
Map<String, ?> sourceOffset = Collections.singletonMap("position", item.getOffset());
records.add(new SourceRecord(
sourcePartition,
sourceOffset,
"target-topic",
Schema.STRING_SCHEMA,
item.getKey(),
Schema.STRING_SCHEMA,
item.getValue()
));
}
return records;
}
@Override
public void stop() {
// Close connections
}
@Override
public String version() {
return "1.0.0";
}
}
public class MySinkTask extends SinkTask {
private ErrantRecordReporter reporter;
@Override
public void start(Map<String, String> props) {
// Initialize connection
try {
reporter = context.errantRecordReporter();
} catch (NoSuchMethodError e) {
reporter = null; // Older Connect runtime
}
}
@Override
public void put(Collection<SinkRecord> records) {
for (SinkRecord record : records) {
try {
writeToDestination(record);
} catch (Exception e) {
if (reporter != null) {
reporter.report(record, e); // Send to DLQ
} else {
throw new ConnectException("Write failed", e);
}
}
}
}
@Override
public void flush(Map<TopicPartition, OffsetAndMetadata> offsets) {
// Ensure all data is persisted before offset commit
}
@Override
public void stop() {
// Close connections
}
@Override
public String version() {
return "1.0.0";
}
}

Register connector class in META-INF/services/:

META-INF/services/org.apache.kafka.connect.source.SourceConnector
com.example.MySourceConnector
# META-INF/services/org.apache.kafka.connect.sink.SinkConnector
com.example.MySinkConnector
@Override
public void start(Map<String, String> props) {
Map<String, Object> partition = Collections.singletonMap("source", connectionUrl);
Map<String, Object> offset = context.offsetStorageReader().offset(partition);
if (offset != null) {
Long lastPosition = (Long) offset.get("position");
seekToPosition(lastPosition);
}
}

Kafka 4.2 includes several Kafka Connect improvements:

FeatureKIPDescription
External schema in JsonConverterKIP-1054schema.content configuration for external schemas, reducing message sizes
Allowlist override policyKIP-1188New ConnectorClientConfigOverridePolicy based on allowlists, enabling fine-grained control over which client configurations connectors may override

Connector Client Configuration Override Policy (Kafka 4.2+)

Section titled “Connector Client Configuration Override Policy (Kafka 4.2+)”

In Kafka 4.2+ (KIP-1188), the new "Allowlist" ConnectorClientConfigOverridePolicy provides a middle ground between None (no overrides permitted) and All (all overrides permitted). Administrators can specify an explicit list of client configuration properties that connectors are allowed to override, restricting access to sensitive settings while still enabling necessary customization.

connector.client.config.override.policy=org.apache.kafka.connect.connector.policy.AllowListConnectorClientConfigOverridePolicy
connector.client.config.override.policy.allowed.configs=batch.size,linger.ms,compression.type