Kafka connector
Reuse a Kafka connection
Use a Kafka Connection credential to share broker and authentication settings between Kafka producer and Kafka consumer connectors.
Reusable credentials require Camunda 8.10 or later and a Kafka element template with the optional Connection credential chooser in the Connection section. Select the same credential on producer tasks and consumer events, or leave Connection credential empty to keep configuring brokers and authentication inline in the Connection section.
A Kafka Connection credential stores these required fields:
| Field | Description |
|---|---|
bootstrapServers | Bootstrap server addresses, separated by commas when you use more than one server. |
username | Kafka username. For a secret reference, use camunda.secrets.MY_KAFKA_USERNAME. |
password | Kafka password. For a secret reference, use camunda.secrets.MY_KAFKA_PASSWORD. |
Create referenced secrets before using the credential. Inside a credential, use camunda.secrets.NAME, not the legacy {{secrets.NAME}} syntax. Legacy secret references remain supported in inline connector fields.
Selecting a credential binds the whole object to kafkaConnectionConfiguration using =camunda.vars.env.<name>, where <name> is the credential's cluster variable name. The credential supplies authentication using SASL_SSL with the PLAIN mechanism. Keep custom authentication configured inline.
The selected credential replaces inline broker and authentication settings. An invalid credential fails execution or consumer activation; the connector doesn't fall back to inline values. Topic, consumer group, offsets, schema settings, and message configuration remain local to each task or event.
Additional properties still overrides the final connection and security properties, including bootstrap.servers. Selecting a credential doesn't restrict the destination or prevent these overrides.
Test a Kafka connection
Test connection performs a read-only Kafka metadata probe to check connectivity and authentication, using the connector runtime's trust configuration. It uses bounded timeouts and reports a failure if the connection can't be established.
The result applies only to the stored credential, not to Additional properties overrides configured on a task or event.
The test doesn't publish or consume messages, join a consumer group, or change offsets. Success doesn't verify topic, consumer group, publishing, or consumption permissions.
Updating a credential follows the existing consumer activation and reload lifecycle. Don't rely on credential edits automatically rotating the connection of an already-running consumer.
- Kafka Producer connector
- Kafka Consumer connector
The Kafka Producer connector is an outbound connector that allows you to connect your BPMN service with Apache Kafka to produce messages.
Prerequisites
To use the Kafka Producer connector, you must have a Kafka instance with a configured bootstrap server.
Use secrets to avoid exposing your sensitive data as plain text. To learn more, see managing secrets.
Create a Kafka Producer connector task
You can apply a connector to a task or event via the append menu. For example:
- From the canvas: Select an element and click the Change element icon to change an existing element, or use the append feature to add a new element to the diagram.
- From the properties panel: Navigate to the Template section and click Select.
- From the side palette: Click the Create element icon.
In each of these menus, you can search by connector name or by the operation you want to perform, such as upload object or send email. Connectors that provide several operations list them as separate entries, and selecting an operation applies the connector with that operation preselected.
After you have applied a connector to your element, follow the configuration steps or see using connectors to learn more.
Make your Kafka Producer connector for publishing messages executable
To make your Kafka Producer connector for publishing messages executable, complete the following sections.
Connection
In the Connection section, select a Connection credential or configure the connection inline:
- Set Bootstrap servers to the URL of the bootstrap server(s). If more than one server is required, use comma-separated values.
- (Optional) Set Username and Password. For example,
{{secrets.MY_KAFKA_USERNAME}}.
Schema
In the Kafka section:
- Select the schema strategy for your messages.
- Select No schema, Inline schema for Avro serialization.
- Select Schema registry if you have a Confluent Schema Registry.
- Set the topic name.
- (Optional) Set producer configuration values in the Headers field. Only
UTF-8strings are supported as header values. - (Optional) Set producer configuration values in the Additional properties field.
The appendix provides more information about:
- Kafka secure authentication.
- Inline schema and Schema registry.
- Pre-configured producer configuration values for this connector.
Additionally, to learn more about supported producer configurations, see the official Kafka documentation.
Message
In the Message section, set the Key and the Value that will be sent to Kafka topic.
Schema strategies
Use Schema strategies with caution, as this is an alpha feature. Functionality may not be comprehensive and could change.
This connector supports different schema strategies, offering a compact, fast, and binary data exchange format for Kafka messages.
When using a schema strategy, each message is serialized according to a specific schema written in JSON format. This schema defines the Kafka message structure, ensuring the data conforms to a predefined format, and enables schema evolution strategies.
To learn more about Schema strategies, refer to the official documentation:
- Inline Avro serialization and official Apache Avro documentation.
- Confluent Schema Registry (Avro, and JSON schemas).
No schema
Select No schema to send messages without a schema. This option is suitable for simple messages that do not require a schema.
Inline schema
Select Inline schema to send messages with an Avro schema.
- This option is suitable for messages that require a schema and that are not (or do not need to be) registered in a schema registry.
- Enter the Avro schema that defines the message structure into the Schema field that appears in the Message section.
Schema registry
Select Schema registry to send messages with a schema registered in a schema registry.
- This option is suitable for messages that require a schema and that are registered in a schema registry.
- You must provide:
- The schema registry URL in the Kafka section.
- The schema itself (that defines the message structure) in the Message section.
- The credentials for the schema registry (if required). Refer to the Schema Registry documentation for more information.
Currently, the Kafka connector supports only Confluent Schema Registry. Other schema registry implementations are not supported at this time.
Example Avro schema and data
The following is an example Avro schema and data:
Avro schema:
{
"doc": "Sample schema to help you get started.",
"fields": [
{
"name": "name",
"type": "string"
},
{
"name": "age",
"type": "int"
},
{
"name": "emails",
"type": {
"items": "string",
"type": "array"
}
}
],
"name": "sampleRecord",
"namespace": "com.mycorp.mynamespace",
"type": "record"
}
Kafka message
-
Key:
employee1 -
Value:
{
"name": "John Doe",
"age": 29,
"emails": ["johndoe@example.com"]
}
Kafka Producer connector response
The Kafka Producer connector returns metadata for a record that has been acknowledged by the Kafka instance.
The following fields are available in the response variable:
timestamp: The timestamp of the message.offset: The message offset.partition: The message partition.topic: The topic name.
For more information on these fields, refer to the official Kafka documentation.
You can use an output mapping to map the response:
-
Use Result Variable to store the response in a process variable. For example,
myResultVariable. -
Use Result Expression to map fields from the response into process variables. For example:
= {
"messageAcknowledgedAt": response.timestamp
}
Appendix and FAQ
What mechanism is used to authenticate against Kafka?
If the fields Username and Password are not empty, by default the Kafka Producer connector enables the credentials-based SASL SSL authentication and the following properties are set:
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='<Your Username>' password='<Your Password>';
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
If any of the fields are not populated, you must configure your security method for your Kafka configuration. You can do this using the Additional properties field.
What are default Kafka Producer client properties?
-
Authentication properties (only if both Username and Password are not empty):
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='<Your Username>' password='<Your Password>';
security.protocol=SASL_SSL
sasl.mechanism=PLAIN -
Bootstrap server property:
bootstrap.servers=<bootstrap server(s) from BPMN> -
Message properties:
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer -
Miscellaneous properties:
session.timeout.ms=45000
client.dns.lookup=use_all_dns_ips
acks=all
delivery.timeout.ms=45000
What is the precedence of client properties loading?
Properties loading consists of three steps:
- Construct client properties from the BPMN diagram: authentication, bootstrap server, message properties. If selected, the Connection credential supplies authentication and bootstrap servers instead of the inline fields.
- Load miscellaneous properties.
- Load and override properties from the field Additional properties.
How do I set or override additional client properties?
The following example sets a new client property client.id and overrides the SASL mechanism to SCRAM SHA-256 instead of plain text:
= {
"client.id":"MyDemoClientId",
"sasl.mechanism":"SCRAM-SHA-256"
}
The Kafka Consumer connector allows you to consume messages by subscribing to Kafka topics and mapping them to your BPMN processes as start or intermediate events.
Prerequisites
To use the Kafka Consumer connector, you must have a Kafka instance with a configured bootstrap server.
Use secrets to avoid exposing your sensitive data as plain text. To learn more, see managing secrets.
Create a Kafka Consumer connector event
- Add a Start Event or an Intermediate Event to your BPMN diagram to get started.
- Change its template to a Kafka Consumer.
- Fill in all required properties.
- Complete your BPMN diagram.
- Deploy the diagram to activate the Kafka consumer.
Configure your Kafka Consumer connector
To make your Kafka Consumer connector executable, fill in the required properties.
Connection
In the Connection section, select a Connection credential or configure the connection inline:
- Set Bootstrap servers to the URL of the bootstrap server(s). If more than one server is required, use comma-separated values.
- Select the Authentication type. If you selected Credentials, set Username and Password.
- Use secrets to avoid exposing your sensitive data as plain text. To learn more, see managing secrets.
- To learn more about Kafka authentication, see Kafka secure authentication.
Kafka properties
In the Kafka section, you can configure the following properties:
- Consumer Group ID: Set the consumer group ID for this connector. Always provide an explicit, stable value that identifies the logical consumer group (for example,
my-app-order-processor). If you leave this field empty, the connector auto-generates an ID from its internal deduplication key. That generated ID can change across connector upgrades, including from 8.8 to 8.9, causing Kafka to treat the connector as a new consumer group and potentially replay already processed messages. - Schema strategy: Select the schema strategy for your messages.
- Select No schema, Inline schema for Avro serialization.
- Select Schema registry If you have a Confluent Schema Registry.
- Topic: Set the topic name.
- Additional properties: Set consumer configuration values.
- Offsets: Set the offsets for the partition. The number of offsets specified should match the number of partitions on the current topic.
- Auto offset reset: Set the strategy to use when there is no initial offset in Kafka or if the specified offsets do not exist on the server.
The appendix provides more information about pre-configured consumer configuration values for this connector.
Additionally, to learn more about supported consumer configurations, see the official Kafka documentation.
Modify an existing inbound Kafka connector
Editing the Consumer Group ID or Offsets properties on an already deployed inbound Kafka connector element can cause unexpected behavior. Delete the element and create a new one instead of editing these properties in place.
- Consumer Group ID: If you clear this field on an existing element, the connector generates a new ID from its internal deduplication key instead of reusing the previous one. Kafka then treats the connector as a new consumer group, which can cause it to replay already processed messages.
- Offsets: If the number of offsets you provide does not match the topic's actual partition count, activation fails with an error.
This is one example of a general pattern for inbound connectors. See modify an existing inbound connector element for more details.
Schema strategy
No schema
Select No schema to send messages without a schema. This option is suitable for simple messages that don’t require a schema.
Inline schema
Select Inline schema to send messages with an Avro schema.
- This option is appropriate for messages that require a schema but are not (or do not need to be) registered in a schema registry.
- Enter the Avro schema that defines the message structure into the Schema field in the Message section.
Schema registry
Select Schema registry to send messages using a schema registered in a schema registry.
- This option is appropriate for messages that require a schema and are registered in a schema registry.
- You must provide:
- The schema registry URL in the Kafka section.
- The schema itself (defining the message structure) in the Message section.
- The credentials for the schema registry, if required. See the Schema Registry documentation for more details.
Currently, the Kafka connector supports only the Confluent Schema Registry. Other schema registry implementations are not supported at this time.
Schema configuration is required only for the outbound connector. It is not required when using Inbound Connectors.
Example Avro schema and data
If the expected Kafka message looks like this:
-
Key:
employee1 -
Value:
{
"name": "John Doe",
"age": 29,
"emails": ["johndoe@example.com"]
}
The corresponding Avro schema to describe this message's structure would be:
{
"doc": "Sample schema to help you get started.",
"fields": [
{
"name": "name",
"type": "string"
},
{
"name": "age",
"type": "int"
},
{
"name": "emails",
"type": {
"items": "string",
"type": "array"
}
}
],
"name": "sampleRecord",
"namespace": "com.mycorp.mynamespace",
"type": "record"
}
This schema defines a structure for a record that includes a name (string), an age (integer), and emails (an array of strings), aligning with the given Kafka message's value format.
Activation condition
Activation condition is an optional FEEL expression field that allows for the fine-tuning of the connector activation. This condition filters if the process step triggers when a Kafka message is consumed.
For example, =(value.itemId = "a4f6j2") only triggers the start event or continues the catch event if the Kafka message has a matching itemId in the incoming message payload. Leave this field empty to trigger your process every time.
By default, this connector does not commit the offset if the message cannot be processed. This includes cases where the activation condition is not met. This means that if there is a message in the topic that cannot be processed due to an activation condition mismatch, the Kafka subscription will be stopped.
Follow the steps below to configure this behavior.
To ignore messages that do not meet the activation condition and commit the offset, select the Consume unmatched events checkbox.
| Consume unmatched events checkbox | Activation condition | Outcome |
|---|---|---|
| Checked | Matched | connector is triggered, offsets are commited |
| Unchecked | Matched | connector is triggered, offsets are commited |
| Checked | Unmatched | connector is not triggered, offsets are commited |
| Unchecked | Unmatched | connector is not triggered, offsets are not commited |
Upgrade from a version without the Consume unmatched events checkbox
If your inbound Kafka connector element was deployed before the Consume unmatched events checkbox existed, the underlying property is absent from that element and defaults to unchecked. Upgrading the runtime does not change this default, so the element keeps its original behavior: it does not commit the offset when a message does not match the activation condition.
To adopt the checked behavior on an existing element, update its element template and redeploy it. Creating a new element with a current template also picks up the new default of checked.
This is one example of a general pattern for inbound connectors. See modify an existing inbound connector element for more details.
Correlation
The Correlation section allows you to configure the message correlation parameters.
The Correlation section is not applicable for the plain start event element template of the Kafka connector. Plain start events are triggered by process instance creation and do not rely on message correlation.
Correlation key
- Correlation key (process) is a FEEL expression that defines the correlation key for the subscription. This corresponds to the Correlation key property of a regular message intermediate catch event.
- Correlation key (payload) is a FEEL expression used to extract the correlation key from the incoming message. This expression is evaluated in the connector Runtime and the result is used to correlate the message.
For example, given that your correlation key is defined with myCorrelationKey process variable, and the incoming Kafka message contains value:{correlationKey:myValue}, your correlation key settings would be as follows:
- Correlation key (process):
=myCorrelationKey - Correlation key (payload):
=value.correlationKey
You can also use the key of the message to accomplish this in the Correlation key (payload) field with =key.
To learn more about correlation keys, see messages.
Message ID expression
The Message ID expression is an optional field that allows you to extract the message ID from the incoming message. The message ID serves as a unique identifier for the message and is used for message correlation. This expression is evaluated in the connector Runtime and the result is used to correlate the message.
In most cases, it is not necessary to configure the Message ID expression. However, it is useful if you want to ensure message deduplication or achieve a certain message correlation behavior.
To learn more about how message IDs influence message correlation, see messages.
For example, if you want to set the message ID to the value of the transactionId field in the incoming message, you can configure the Message ID expression as follows:
= value.transactionId
Message TTL
The Message TTL is an optional field that allows you to set the time-to-live (TTL) for the correlated messages. TTL defines the time for which the message is buffered in Zeebe before being correlated to the process instance (if it can't be correlated immediately).
The value is specified as an ISO 8601 duration. For example, PT1H sets the TTL to one hour. Learn more about the TTL concept in Zeebe in the message correlation guide.
Deduplication
The Deduplication section allows you to configure the connector deduplication parameters.
Connector deduplication is a mechanism in the connector Runtime that determines how many Kafka subscriptions are created if there are multiple occurrences of the Kafka Consumer connector in the BPMN diagram. This is not to be confused with message deduplication.
By default, the connector runtime deduplicates connectors based on properties, so that elements with the same subscription properties only result in one subscription.
To learn more about deduplication, see deduplication.
To customize the deduplication behavior, select the Manual mode checkbox, and configure the custom deduplication ID.
Output mapping
The Kafka Consumer connector returns the consumed message.
The following fields are available in the response variable:
key: The key of the message.value: The value of the message.rawValue: The value of the message as a JSON string.
You can use an output mapping to map the response:
- Use Result variable to store the response in a process variable. For example,
myResultVariable. - Use Result expression to map fields from the response into process variables. For example:
= {
"itemId": value.itemId
}
Activate the Kafka Consumer connector by deploying your diagram
When you click the Deploy button, your Kafka Consumer is activated and starts consuming messages from the specified topic.
Appendix and FAQ
What mechanism is used to authenticate against Kafka?
If you selected Credentials as the Authentication type and the fields Username and Password are not empty, by default the Kafka Consumer connector enables the credentials-based SASL SSL authentication, and sets the following properties:
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='<Your Username>' password='<Your Password>';
security.protocol=SASL_SSL
sasl.mechanism=PLAIN
If any of the field is not populated, you must configure your security method for your Kafka configuration. You can do this using the Additional properties field.
What are default Kafka Consumer client properties?
-
Authentication properties (only if both Username and Password are not empty):
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='<Your Username>' password='<Your Password>';
security.protocol=SASL_SSL
sasl.mechanism=PLAIN -
Bootstrap server property:
bootstrap.servers=<bootstrap server(s) from BPMN> -
Message properties:
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer -
Miscellaneous properties:
session.timeout.ms=45000
client.dns.lookup=use_all_dns_ips
acks=all
group.id=kafka-inbound-connector-{{bpmnProcessId}}
enable.auto.commit=false
The group.id value above is auto-generated when no explicit Consumer Group ID is configured in the connector. This generated ID is derived from the connector's internal deduplication key and can change across connector upgrades, including from 8.8 to 8.9. When the group ID changes, Kafka treats the connector as a new consumer group, which means committed offsets are not reused and messages may be replayed. To avoid this, always set an explicit Consumer Group ID. You can look up existing consumer groups to find the current group ID in use.
What is the precedence of client properties loading?
Properties loading consists of three steps:
- Construct client properties from the BPMN diagram: authentication, bootstrap server, message properties. If selected, the Connection credential supplies authentication and bootstrap servers instead of the inline fields.
- Load miscellaneous properties.
- Load and override properties from the field Additional properties.
How is the message payload deserialized?
As Kafka messages usually use JSON format, we first try to deserialize it as a JsonElement. If this fails (for example, because of a wrong format) we use the String representation of the original raw value. For convenience, we always store the original raw value as String in a different attribute.
The deserialized object structure:
{
key: "String"
rawValue: "String"
value: {}
}
When is the offset committed? What happens if the connector execution fails?
The following outcomes are possible:
- If the connector execution is successful and the Activation condition was met, the offset is committed.
- If the Activation condition was not met, the offset is also committed to prevent consuming the same message twice.
- If the connector execution fails due to an unexpected error (for example, Zeebe is unavailable), the offset is not committed.
What lifecycle does the Kafka Consumer connector have?
The Kafka Consumer connector is a long-running connector that is activated when the process is deployed, and deactivated when the process is undeployed or overwritten by a new version.