Read Apache Kafka topics into FrameworX tags.
Connectors Library → Browse By Connector Group → IT-Cloud and AI Connectors → Kafka Consumer Connector
Available from FrameworX 10.1.6. Earlier versions do not include this connector.
- Name: KafkaConsumer
- Version 1.0.0.0
- Protocol: Kafka (consumer)
- Interface: TCP/IP
- Runtime: Windows x64, .NET Framework 4.8 and .NET 10
- Available from: FrameworX 10.1.6
- Configuration:
- Devices / Protocols
Overview
The Kafka Consumer connector reads messages from an Apache Kafka cluster into FrameworX tags. Each point subscribes to a topic and takes the latest matching message: the whole message, one field of a JSON message, or the value, quality and timestamp of a message from the Kafka Producer Connector. Each tag holds the latest value of its topic, not a log of every message. The connector does not create tags from topics.
This page covers the channel, node and point settings, the message mapping and delivery behavior. To publish tag values to Kafka, see Kafka Producer Connector. For MQTT brokers, see MQTT Client Connector. For general channel, node and point concepts, see Devices Module.
Communication Driver Information | |
|---|---|
Driver name | KafkaConsumer |
Assembly Name | T.ProtocolDriver.KafkaConsumer |
Assembly Version | 1.0.0.0 |
Protocol description | Kafka Consumer Connector |
Available for Linux | False |
Direction | Ingress only (Kafka to FrameworX) |
Prerequisites
- FrameworX 10.1.6 or later.
- A Windows x64 computer running the FrameworX runtime. See Platform Support below for Windows ARM64.
- The Microsoft Visual C++ 2015-2022 x64 runtime. The FrameworX installer installs it by default. If you clear this component in a Custom install and the computer does not have the runtime, the connector reads no messages.
- An Apache Kafka broker reachable from the runtime computer on its listener port, for example
9092. - The topics to read exist on the broker.
- For SASL authentication or a cluster with access control: a Kafka user with read permission on the topics and on the consumer group of the node. See Consumer Group below for the group name.
Procedure
- In Devices → Channels, create a channel and select the protocol Kafka Consumer Connector.
- Open the channel Protocol Options and set the consumer group, start position, payload format and timestamp source. See Channel Configuration.
- In Devices → Nodes, create a node on the channel and set the Station fields: broker list, security protocol and credentials. See Node Configuration.
- In Devices → Points, add one point per tag to update. Set Address to
Topic;Key;ValuePathand AccessType to a read type with polling, for exampleRead. See Points Configuration. - Start the runtime. On each polling cycle, each tag receives the latest matching message.
Channel Configuration
Protocol Options
Field order in the options string: GroupId;AutoOffsetReset;PayloadFormat;TimestampSource
Option | Values | Default | Description |
|---|---|---|---|
GroupId | Text, or blank | (blank) | Kafka consumer group. See Consumer Group. |
AutoOffsetReset |
|
|
|
PayloadFormat |
|
|
|
TimestampSource |
|
|
|
The runtime reads the Protocol Options at start. After a change, restart the runtime.
Options string examples:
Default values: ;Latest;Auto;Payload Restore the last value of each key at start: ;Earliest;Auto;Payload Named group, text payload, arrival time: PLANT1.SCADA;Latest;Text;Arrival
Consumer Group
Each node consumes as its own Kafka consumer group. With Latest, the group keeps the read position of the node, so a restart resumes where the node stopped.
- GroupId blank: the group name is
FrameworX.<ComputerName>.<SolutionName>.<CHANNEL>.<NODE>, with the channel and node names in upper case. Two solutions, or two computers, reading the same topic therefore use different groups, and each receives every message. - GroupId set: the group name is
<GROUPID>.<NODE>. The runtime converts the GroupId to upper case, soplant1.scadabecomesPLANT1.SCADA. Do not use=or;in the GroupId. Grant the Kafka read permission to the upper-case group name. - Two runtimes of the same solution on the same computer, for example a Development and a Production profile, share the default group and split the topic partitions between them. Set a different GroupId on one of them.
Node Configuration
Station Configuration
Station syntax: <BootstrapServers>;<SecurityProtocol>;<SaslMechanism>;[SaslUsername];[SaslPassword];[SslCaLocation]
The Station fields are the same as in the Kafka Producer Connector, so a producer node and a consumer node of the same cluster use the same Station string. The node uses its primary Station only. The connector does not switch to a backup Station.
Field | Values | Default | Description |
|---|---|---|---|
BootstrapServers |
|
| Comma-separated list of brokers, for example |
SecurityProtocol |
|
| Security of the broker connection. |
SaslMechanism |
|
| SASL mechanism. Used only when SecurityProtocol is |
SaslUsername | Text | (blank) | SASL user name. Used only with |
SaslPassword | Password | (blank) | SASL password. The Designer stores it encoded. Used only with |
SslCaLocation | File path | (blank) | Path to a PEM file with the CA certificate, or certificate chain, of the broker, for example |
Enter SaslPassword in the Designer station editor. The station string separates fields with semicolons, and the Designer encoding keeps a password with a semicolon in one field. A plain-text password with a semicolon, or one starting with #, written to the station string by an import or a script, reaches the broker altered.
Points Configuration
Address
Address syntax: <Topic>;[Key];[ValuePath]
Field | Default | Description |
|---|---|---|
Topic | (required) | Kafka topic to read. Letters, digits, |
Key | (blank) | Blank: every message of the topic. Set: only messages with this key. A message without a key matches when its JSON |
ValuePath | (blank) | Blank: the mapping rules in Message Mapping. Set: a JSONPath into the message, for example |
Address examples:
Messages with key Line1Temp on topic plant.line1: plant.line1;Line1Temp; Every message on topic sensors.raw: sensors.raw;; Field temp of the JSON messages on sensors.json: sensors.json;;$.temp
AccessType
Assign a read AccessType with polling, for example the predefined Read type. The polling period is the tag update period: on each polling cycle the tag receives the latest matching message received since start. The connector ignores writes to its points. See Devices AccessTypes Reference.
Message Mapping
For each point, the connector maps a matching message to a value, a quality and a timestamp, in this order:
- PayloadFormat
Text: the whole message as a string, quality Good (192). The ValuePath is not used. - ValuePath set: the JSON value at the ValuePath, quality Good. A message without a value at this path, or not in JSON format, leaves the tag unchanged.
- Message from the Kafka Producer Connector: a JSON object with
tagandvaluefields. The tag receivesvalue, thequalityfield when it is an integer, and thetimestampfield with TimestampSourcePayload. - Any other message: the whole message as a string, quality Good.
Producer message example:
{"tag":"K_PlainCounter","value":265,"quality":192,"timestamp":"2026-10-08T10:11:00.6674091Z"}
- Types: JSON numbers become integer or floating-point values,
trueandfalsebecome Boolean values, strings stay strings, and JSON objects and arrays reach the tag as JSON text. The runtime converts the value to the tag type. - Not mapped: a JSON
nullvalue, and a Kafka record without a payload, for example a delete marker in a compacted topic, leave the tag unchanged. - Dates in JSON strings stay strings.
- Quality: the connector uses the
qualityfield of Kafka Producer Connector messages only. For other messages it sets Good.
Delivery Behavior
- Latest value per tag. Between two polling cycles, the last matching message wins. Intermediate values between two cycles do not reach the tag or the historian. A shorter polling period skips fewer values. To record every message, read the topic with another Kafka client.
- Message order. Kafka orders messages within a partition. When the messages of one point come from several partitions, for example a blank Key on a multi-partition topic, the tag keeps the message with the newest Kafka record time.
- Catch-up at start. When the node starts, or Kafka reassigns its partitions, the connector reads up to the end of each partition before it updates the tags. A replay with
Earliesttherefore writes only the last value of each point, not the intermediate values. A partition without data or end-of-partition for 10 seconds counts as read. - Startup. Until Kafka assigns the partitions to the node, the tags keep their startup value and quality. When no assignment arrives within 30 seconds, the tags go Bad after the channel retries.
- Connection loss. When the client loses all brokers, or the broker rejects the credentials, the tags go Bad after the channel retries. While the brokers are unreachable, the connector checks the cluster every 2 seconds. When the connection returns, the tags go Good again with the recovery time as timestamp, so the historian keeps the Bad sample before the Good one. A point with no message since start stays Bad until its first message.
- Read position. The consumer group commits the read position on its own. A message read but not yet applied to a tag when the runtime stops is not read again with
Latest. - Timestamps of replayed values. A value replayed at start, or received after a connection loss, carries the start or recovery time when its own timestamp is older. This applies to every TimestampSource.
- Configuration changes. When the node Station changes, the connector rebuilds the consumer with the new settings. A change of BootstrapServers also discards the values received from the previous cluster. Each tag keeps its current value until a message from the new cluster arrives.
- Loops. Do not read and publish the same tag on the same topic.
Station Examples
1. Local Apache Kafka, no authentication, no TLS
Field | Value |
|---|---|
BootstrapServers |
|
SecurityProtocol |
|
SaslMechanism |
|
SaslUsername | (blank) |
SaslPassword | (blank) |
SslCaLocation | (blank) |
Station string:
localhost:9092;Plaintext;Plain;;;
Channel Protocol Options: the default values. Point: Address sensors.raw;;, AccessType Read.
2. Apache Kafka with SASL_SSL, SCRAM-SHA-256 and a private CA file
Field | Value |
|---|---|
BootstrapServers |
|
SecurityProtocol |
|
SaslMechanism |
|
SaslUsername |
|
SaslPassword | The password of |
SslCaLocation |
|
The broker listener on port 9093 uses SASL_SSL with SCRAM-SHA-256, and ca.pem holds the CA certificate in PEM format. The broker certificate lists kafka01.plant.local as a host name. Create the SCRAM user on the broker with kafka-configs, and grant it read permission on the topics and on the consumer group, before you start the runtime.
Platform Support
Platform | Supported |
|---|---|
Windows x64, .NET Framework 4.8 runtime | Yes |
Windows x64, .NET 10 runtime | Yes |
Windows ARM64, .NET Framework 4.8 runtime | Yes, under x64 emulation |
Windows ARM64, .NET 10 runtime started with the x64 | Yes, under x64 emulation |
Windows ARM64, .NET 10 runtime with the default ARM64 | No |
Linux, any architecture | No |
The connector depends on the native Kafka client library for Windows x64. On a host where this library does not load, the node stays Bad.
Troubleshooting
The connector writes these messages to the Trace Window (see Runtime Diagnostics Reference). Check Devices in the Trace Window settings. Messages marked (Debug) appear only when you also check Debug. For rows with a quoted message, the first column shows the start of the message.
Symptom | Likely Cause | Resolution |
|---|---|---|
| The first Station field is blank. | Enter at least one broker as |
| The topic name has a character Kafka does not accept, or the ValuePath is not a valid JSONPath. | Correct the point Address. See Address. |
| No broker reachable: broker stopped, wrong host or port, or TLS settings different from the broker listener. | Confirm the runtime computer reaches the broker port. Match SecurityProtocol to the listener on the broker port. Tags with a received message return to Good when the cluster is reachable again, with no runtime restart. |
| Wrong SASL user name, password or mechanism. | Check SaslUsername, SaslPassword and SaslMechanism against the broker user. |
| Kafka did not assign partitions to the node. Common causes: the topics do not exist, the user has no read permission on the consumer group, or the group waits for a previous runtime instance after a forced stop. | Create the topics. Grant read permission on the consumer group name, in upper case when you set a GroupId (see Consumer Group). After a forced stop of the runtime, Kafka releases the previous group member within 45 seconds by default. After the assignment, each tag returns to Good on its next message, or at start with |
| The client rejected the configuration, for example an SslCaLocation file missing or unreadable. Or the native Kafka client library did not load, see the Windows ARM64 and Visual C++ rows. | Check the reason at the end of the message. The connector retries every 5 seconds, so a corrected Station or CA file takes effect without a restart. |
| The client reported an error for a topic, for example a subscribed topic missing on the broker, or a connection error to one broker. The client keeps retrying. | Read the reason at the end of the message. Create the missing topic, or check the broker named in the message. |
| A partition sent neither data nor an end-of-partition signal for 10 seconds after an assignment. | No action when the tags update. Check the broker health when the message repeats. |
| The point has a ValuePath, and the message is not JSON or has no value at this path. Topics carrying several message layouts produce this message for the points of the other layouts. | No action for mixed topics. Otherwise check the ValuePath against a message read with |
A tag keeps its startup value after a restart, no error | AutoOffsetReset is | Expected with |
A tag never updates, no error | The Key does not match the message key, including case, or the point AccessType has no read polling. | Read the message keys with |
TLS handshake fails with | The broker certificate does not list the host name used in BootstrapServers. The client checks the host name of the broker certificate. | Use the host name from the broker certificate in BootstrapServers, or issue a broker certificate with this host name in its Subject Alternative Name. |
TLS handshake fails with a broker certificate trusted by other clients | The TLS library rejects certificates with RSA keys shorter than 2048 bits. It also does not read a system | Use a broker certificate with an RSA key of 2048 bits or more. Set the CA file in SslCaLocation instead of an OpenSSL configuration file. |
| The default ARM64 | Use the .NET Framework 4.8 runtime, or start the .NET 10 runtime with the x64 |
| The Visual C++ 2015-2022 x64 runtime is missing. | Install the Microsoft Visual C++ 2015-2022 Redistributable (x64). |
Third-Party Components
The connector uses the same components as the Kafka Producer Connector, installed in the Protocols folder of the FrameworX installation:
Component | Version |
|---|---|
Confluent.Kafka (.NET client) | 2.15.1 |
librdkafka | 2.15.1 |
OpenSSL | 3.5.9 LTS |
libcurl | 8.21.0 |
zlib | 1.3.2 |
Zstandard (zstd) | 1.5.7 |
License notices for these components ship in <FrameworX installation>\Protocols\Kafka-ThirdPartyNotices.txt.
Driver Revision History
KafkaConsumer Revision History | |
|---|---|
Version | Notes |
1.0.0.0 | Initial release in FrameworX 10.1.6. |
In this section...