skip to content

Single Message Transforms and Predicates

Single Message Transforms and predicates for light per-record reshaping and routing inside Connect. Interviewers ask where SMTs are enough and where you should reach for a real stream processor.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

What is a Single Message Transform (SMT) in Kafka Connect, and how do you configure a chain of them on a connector?

level: juniorimportance: must knowfreq 70%

answer

  1. transforms=<alias list>, runs left-to-right
  2. transforms.<alias>.type = FQ class
  3. per-record, stateless
  4. $Key vs $Value inner classes
  5. implements Transformation<R>

basics

~20 s

An SMT is a small function that modifies each record as it flows through a Kafka Connect connector. You list transforms by alias in the transforms config, then configure each alias with transforms.<alias>.type and its properties. They run in the order listed.

solid answer

~40 s

A Single Message Transform (SMT) is a lightweight, stateless function applied to each individual record inside a Kafka Connect connector pipeline. For a source connector SMTs run after the connector produces records but before they hit Kafka; for a sink connector they run after reading from Kafka but before delivery to the sink. You declare a chain with the `transforms` property, listing comma-separated aliases (e.g. `transforms=route,addTs`). Each alias is then configured with `transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter` plus that SMT's own properties like `transforms.route.regex` and `transforms.route.replacement`. The chain executes left-to-right in the listed order, each SMT receiving the previous one's output. SMTs implement the `org.apache.kafka.connect.transforms.Transformation` interface and are meant for simple per-record edits — not joins, aggregations, or anything needing state across records.

go deeper

for a junior

Know SMTs edit one record at a time and that transforms= lists aliases configured by transforms.<alias>.type.

for a middle

Explain source-vs-sink ordering relative to the converter, the $Key/$Value variants, and chain execution order.

for a senior

Discuss schema propagation, statelessness limits, and when to reach for Kafka Streams instead.

for a principal

Frame SMTs as a thin, deterministic hot-path concern; set org guidance on what belongs in SMTs vs stream processing vs upstream connectors.

## What problem SMTs solve Kafka Connect moves data between Kafka and external systems using **source connectors** (external system → Kafka) and **sink connectors** (Kafka → external system). Often the record needs a small adjustment in transit: rename a field, drop a field, change the destination topic, add a timestamp. Rewriting the connector for each tweak would be wasteful. A **Single Message Transform (SMT)** is a reusable, configuration-driven function that operates on **one record at a time** to make exactly these small edits. ## Where in the pipeline they run - **Source connector:** connector task produces a `SourceRecord` → SMT chain → converter serializes → record written to Kafka. - **Sink connector:** record read from Kafka → converter deserializes → SMT chain → `SinkRecord` handed to the sink task. So SMTs always sit between the connector and the converter/Kafka boundary. ## Configuring a chain Three pieces: 1. `transforms` — a comma-separated list of **aliases** you invent, e.g. `transforms=insertTs,route`. 2. `transforms.<alias>.type` — the fully-qualified class name of the SMT implementation. 3. `transforms.<alias>.<property>` — that SMT's own settings. Example: ``` transforms=insertTs,route transforms.insertTs.type=org.apache.kafka.connect.transforms.InsertField$Value transforms.insertTs.timestamp.field=ingestedAt transforms.route.type=org.apache.kafka.connect.transforms.RegexRouter transforms.route.regex=(.*) transforms.route.replacement=prod_$1 ``` The chain runs **in the order the aliases appear** in `transforms`. Here every record first gets an `ingestedAt` field inserted, then its topic is renamed with a `prod_` prefix. ## Key properties of SMTs - **Per-record and stateless:** they see one record and cannot aggregate or join across records. - **Key vs Value variants:** many SMTs ship as two inner classes, `...$Key` and `...$Value`, choosing whether to operate on the record key or value. - **Schema-aware:** SMTs work on both schemaful (e.g. Avro `Struct`) and schemaless (plain `Map`) records, transforming the schema as well as the data when present. - **Implement `Transformation<R>`:** the contract has `apply(R record)`, `configure(Map)`, `config()`, and `close()`. ## What SMTs are NOT for Anything requiring external lookups, buffering, windowing, or cross-record state. For that, use Kafka Streams or ksqlDB. SMTs are intentionally simple to keep them cheap and predictable on the hot path.

  • Do SMTs run before or after the converter on a source connector?
    Before. On a source connector the SMT chain transforms the SourceRecord, then the converter serializes it to bytes for Kafka. On a sink connector the order is reversed: converter deserializes first, then the SMT chain runs.
  • Why can't you use an SMT to join two records or compute a running total?
    SMTs are stateless and operate on a single record at a time — the apply() method only sees the current record. Cross-record operations need a stream processor like Kafka Streams or ksqlDB.

saying these in an interview costs you the question

  • Saying SMTs can aggregate or join records — they are strictly per-record and stateless.
  • Thinking the chain order doesn't matter — it executes strictly in the order aliases are listed.
  • Confusing where SMTs run relative to the converter (it differs between source and sink).

context

open as a page

Walk through the common built-in SMTs (InsertField, ReplaceField, MaskField, RegexRouter, TimestampRouter, ExtractField, Cast, Flatten) and what each does.

level: middleimportance: must knowfreq 65%

basics

~20 s

InsertField adds a field; ReplaceField renames/drops fields; MaskField hides values; RegexRouter and TimestampRouter rewrite the target topic name; ExtractField pulls one field up to be the whole key/value; Cast changes field types; Flatten collapses nested structs into dotted top-level fields.

open as a page

How do predicates work with SMTs in Kafka Connect, and what do the built-in TopicNameMatches and HasHeaderKey predicates plus `negate` let you do?

level: seniorimportance: must knowfreq 50%

basics

~20 s

A predicate is a per-record boolean test you attach to an SMT so the SMT only runs when the predicate is true. You define predicates under predicates, reference one via transforms.<alias>.predicate, and flip it with transforms.<alias>.negate=true. Built-ins include TopicNameMatches (regex on topic) and HasHeaderKey.

open as a page

How do you implement a custom SMT by implementing the Transformation interface, and what are the key methods and gotchas (schema handling, key/value variants, deployment)?

level: seniorimportance: should knowfreq 35%

basics

~20 s

Implement org.apache.kafka.connect.transforms.Transformation<R extends ConnectRecord<R>>. Override apply(R) to return a transformed record, configure(Map) to read config, config() to declare a ConfigDef, and close() to release resources. Handle both schema and schemaless records, build a new record with record.newRecord(...), package it as a plugin, and put the Jar on the Connect plugin path.

open as a page

When designing an SMT chain (e.g. ExtractField + Cast + Flatten + RegexRouter with predicates), why does ordering matter and what are common pitfalls?

level: principalimportance: should knowfreq 25%

basics

~20 s

SMTs run in listed order and each consumes the previous one's output, so an SMT that renames or removes a field changes what later SMTs can see. You must order field-shaping before routing or extraction that depends on those fields, mind schema changes, predicate scoping, and tombstone handling.

open as a page