Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.7 Single Message Transformations (SMTs)

The division of labor, restated:

SMTsKafka Streams
Stateless. Done inside Kafka Connect while copying, “often without writing any code.”Stateful. “More complex transformations, which typically involve joins or aggregation.”

The built-in SMTs in Apache Kafka

SMTWhat it does
CastChange data type of a field
MaskField"Replace the contents of a field with null. Useful for removing sensitive or personally identifying data."
Filter"Drop or include all messages matching a condition. Built-in conditions include matching on a topic name, a particular header, or whether the message is a TOMBSTONE (that is, has a null value)."
Flatten"Transform a nested data structure to a flat one... by concatenating all the names of all fields in the path to a specific value."
HeaderFrom"Move or copy fields from the message into the header."
InsertHeader"Add a static string to the header of each message."
InsertField"Add a new field... either using values from its metadata such as offset, or with a static value."
RegexRouter"Change the destination topic using a regular expression and a replacement string."
ReplaceField"Remove or rename a field."
TimestampConverter"Modify the time format of a field — for example, from Unix Epoch to a String."
TimestampRouter"Modify the topic based on the message timestamp. Mostly useful in SINK connectors when we want to copy messages to specific TABLE PARTITIONS based on their timestamp and the topic field is used to find an equivalent dataset in the destination system."

Third-party collections: GitHub — Lenses.io, Aiven, Jeremy Custenborder — or Confluent Hub. Learning: the "Twelve Days of SMT" blog series; tutorials for writing your own.

Worked example: add a lineage header

Motivation: "The header will indicate that the record was created by this MySQL connector, which is useful in case auditors want to examine the LINEAGE of these records."

echo '{
  "name": "mysql-login-connector",
  "config": {
    "connector.class": "JdbcSourceConnector",
    "connection.url": "jdbc:mysql://127.0.0.1:3306/test?user=root",
    "mode": "timestamp",
    "table.whitelist": "login",
    "validate.non.null": "false",
    "timestamp.column.name": "login_time",
    "topic.prefix": "mysql.",
    "name": "mysql-login-connector",
    "transforms": "InsertHeader",
    "transforms.InsertHeader.type":
      "org.apache.kafka.connect.transforms.InsertHeader",
    "transforms.InsertHeader.header": "MessageSource",
    "transforms.InsertHeader.value.literal": "mysql-login-connector"
  }}' | curl -X POST -d @- http://localhost:8083/connectors \
       --header "content-Type:application/json"

The naming convention to internalize:

transforms                 = <alias>[,<alias2>,...]   ← ordered chain
transforms.<alias>.type    = <fully qualified class>
transforms.<alias>.<param> = <value>

The transforms list is an ordered chain — each alias is then configured by its own type and parameter keys.

Result (needs Apache Kafka 2.7+ to print headers in the console consumer):

bin/kafka-console-consumer.sh --bootstrap-server=localhost:9092 \
  --topic mysql.login --from-beginning --property print.headers=true

NO_HEADERS                             {"schema":...,"payload":{"username":"tpalino",...}}
MessageSource:mysql-login-connector    {"schema":...,"payload":{"username":"rajini",...}}

"the old records show NO_HEADERS, but the new records show MessageSource:mysql-login-connector."

ERROR HANDLING AND DEAD LETTER QUEUES

"Transforms is an example of a connector config that isn't specific to one connector but can be used in the configuration of ANY connector. Another very useful config that can be used in any SINK connector is error.tolerance — you can configure any connector to silently drop corrupt messages, or to route them to a special topic called a 'DEAD LETTER QUEUE.'" → "Kafka Connect Deep Dive — Error Handling and Dead Letter Queues" blog post.

(Compare Ch. 7 §5.2 Rule 5 Pattern B: the same DLQ pattern, but here it's a config flag instead of application code.)


On this page