9.7 Single Message Transformations (SMTs)
The division of labor, restated:
| SMTs | Kafka 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
| SMT | What it does |
|---|---|
| Cast | Change 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 showMessageSource: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.)
9.6 Example: MySQL → Kafka → Elasticsearch
Three options: Confluent Hub client, download from Confluent Hub (or wherever the connector is hosted), or build from source:
9.8 A deeper look at Connect internals
Note: "The JSON converter can be configured to either include a schema in the result record or not — so we can support both structured and semistructured data."