How do predicates work with SMTs in Kafka Connect, and what do the built-in TopicNameMatches and HasHeaderKey predicates plus `negate` let you do?
answer
- KIP-585, Kafka 2.6+
- predicates=<name>; transforms.<alias>.predicate=<name>
- negate=true inverts (run when false)
- TopicNameMatches=pattern, HasHeaderKey=name, RecordIsTombstone
- Filter SMT needs a predicate or it drops everything
basics
~20 sA 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.
solid answer
~40 sPredicates (added in KIP-585, Kafka 2.6) let you apply an SMT **conditionally** per record. You declare them in a `predicates` list, each with `predicates.<name>.type` and its config. Then on a transform you set `transforms.<alias>.predicate=<name>` so the SMT runs only when the predicate returns true, and `transforms.<alias>.negate=true` to invert it (run when false). Built-ins: **TopicNameMatches** (`pattern` = a Java regex) matches the record's topic; **HasHeaderKey** (`name`) is true if a header with that key exists; **RecordIsTombstone** is true for null-value tombstones. A classic pattern is gating a routing or masking SMT to a subset of topics, or using `negate` with a Filter SMT to drop everything that does NOT match. Each SMT takes at most one predicate; combine logic by chaining multiple transforms.
go deeper
Know a predicate is a true/false gate that decides whether an SMT runs on a record.
Wire up predicates, transforms.<alias>.predicate, and negate; name TopicNameMatches and HasHeaderKey.
Use Filter+predicate+negate patterns, handle tombstones, and explain the one-predicate-per-transform limit and KIP-585 origin.
Design conditional-routing topologies, decide when a custom predicate beats chaining, and guard against tombstone NPEs across the fleet.
## Why predicates exist By default every SMT in a chain runs on **every** record. But often you want an SMT to apply only to *some* records — e.g. mask PII only on the `customers` topic, or route only records carrying a particular header. Before KIP-585 you'd need separate connectors. **Predicates** (Apache Kafka 2.6+) add a per-record boolean gate to any single SMT. ## Configuration shape Three config families work together: 1. **Declare predicates** much like transforms: ``` predicates=isCustomers predicates.isCustomers.type=org.apache.kafka.connect.transforms.predicates.TopicNameMatches predicates.isCustomers.pattern=.*customers.* ``` 2. **Attach to a transform** with `predicate`: ``` transforms=mask transforms.mask.type=org.apache.kafka.connect.transforms.MaskField$Value transforms.mask.fields=ssn transforms.mask.predicate=isCustomers ``` Now `mask` runs only on records whose topic matches the pattern. 3. **Invert with `negate`:** ``` transforms.mask.negate=true ``` With `negate=true`, the SMT runs when the predicate is **false** — here, on every topic *except* customers. ## Built-in predicates - **TopicNameMatches** — config `pattern`, a Java `java.util.regex.Pattern`. True when the record's topic matches. Most common predicate. - **HasHeaderKey** — config `name`. True when the record has at least one header with that key. Lets you route/transform based on producer-set headers. - **RecordIsTombstone** — no config. True when the record value is null (a tombstone / delete marker). Useful to skip transforms that would NPE on tombstones, or with `Filter` + `negate` to drop tombstones. ## The Filter SMT partnership `org.apache.kafka.connect.transforms.Filter` drops any record it sees (it has no config). Alone it would drop everything, so it is **almost always paired with a predicate**. To keep only customer records: a Filter gated by `TopicNameMatches` with `negate=true` (drop when NOT customers). To drop tombstones: Filter + `RecordIsTombstone`. ## Semantics and limits - **One predicate per transform.** There is no built-in AND/OR; compose by chaining multiple transforms each with its own predicate, or write a custom predicate. - Predicates are pure per-record tests implementing the `Predicate<R>` interface (`test(R record)`), declared via the `predicates` config and the `Predicate` plugin type. - `negate` is a property of the *transform-predicate binding*, not of the predicate itself — the same predicate can be used negated by one transform and non-negated by another. - Predicates only choose **whether an SMT runs**; they cannot modify the record themselves.
- How would you drop all records that are NOT on the `orders` topic?Use the Filter SMT gated by a TopicNameMatches predicate (`pattern=.*orders.*`) with `negate=true`. Filter drops the record, and negate makes it fire only when the topic does NOT match.
- Can one SMT have two predicates ANDed together?No. Each transform accepts a single `predicate`. To combine conditions, chain multiple transforms each with its own predicate, or implement a custom Predicate that encodes the combined logic.
- Is `negate` a property of the predicate or of the transform?Of the transform's predicate binding (`transforms.<alias>.negate`). The same declared predicate can be applied negated by one transform and un-negated by another.
saying these in an interview costs you the question
- Saying a predicate can modify the record — it only returns true/false to gate an SMT.
- Putting `negate` on the predicate declaration; it belongs on `transforms.<alias>.negate`.
- Assuming you can attach multiple predicates with AND/OR to one transform — only one per transform.
- Using the Filter SMT without a predicate and expecting it to keep records — it drops everything unconditionally.