|
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. |
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) |
9092.Topic;Key;ValuePath and AccessType to a read type with polling, for example Read. See Points Configuration.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 |
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.
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>.<NODE>. The runtime converts the GroupId to upper case, so plant1.scada becomes PLANT1.SCADA. Do not use = or ; in the GroupId. Grant the Kafka read permission to the upper-case group name.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 |
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 |
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.
For each point, the connector maps a matching message to a value, a quality and a timestamp, in this order:
Text: the whole message as a string, quality Good (192). The ValuePath is not used.tag and value fields. The tag receives value, the quality field when it is an integer, and the timestamp field with TimestampSource Payload.Producer message example:
{"tag":"K_PlainCounter","value":265,"quality":192,"timestamp":"2026-10-08T10:11:00.6674091Z"} |
true and false become 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.null value, and a Kafka record without a payload, for example a delete marker in a compacted topic, leave the tag unchanged.quality field of Kafka Producer Connector messages only. For other messages it sets Good.Earliest therefore 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.Latest.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.
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 | 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.
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). |
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.
KafkaConsumer Revision History | |
|---|---|
Version | Notes |
1.0.0.0 | Initial release in FrameworX 10.1.6. |