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)?
answer
- Transformation<R extends ConnectRecord<R>>
- apply / configure / config (ConfigDef) / close
- return null = drop record
- record.newRecord(...) — records immutable
- branch schemaful Struct vs schemaless Map; cache schema
- JAR on plugin.path, classloader isolation
basics
~20 sImplement 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.
solid answer
~40 sA custom SMT implements `Transformation<R extends ConnectRecord<R>>`. Key methods: `configure(Map<String,?>)` parses config (typically into an `AbstractConfig` from your `ConfigDef`); `config()` returns the `ConfigDef` for validation/docs; `apply(R record)` does the work and returns a new record (or null to drop it); `close()` frees resources. Records are immutable, so you produce a transformed copy via `record.newRecord(topic, partition, keySchema, key, valueSchema, value, timestamp[, headers])`. You must handle **two cases**: schemaful records (operate on `Schema` + `Struct`, often caching rebuilt schemas with a `SchemaUpdateCache`) and schemaless records (operate on a `Map`). Conventionally expose `Key` and `Value` static subclasses so users pick which part to transform. Deploy by packaging the class as a plugin JAR and placing it on the worker's `plugin.path`; then reference it by FQ class name in `transforms.<alias>.type`. Keep it stateless, fast, and null-safe (tombstones).
go deeper
Know that a custom SMT implements the Transformation interface and is configured by class name.
Name apply/configure/config/close and that returning null drops a record; understand immutability.
Handle schemaful vs schemaless, schema caching, Key/Value variants, tombstone safety, and plugin.path deployment.
Set conventions for custom-SMT performance (schema cache), classloader isolation, ConfigDef rigor, and when a custom SMT is justified vs composing built-ins.
## The contract A custom SMT implements the generic interface: ``` public interface Transformation<R extends ConnectRecord<R>> extends Configurable, Closeable { R apply(R record); ConfigDef config(); void close(); // configure(Map<String,?>) inherited from Configurable } ``` The type parameter `R` is `SourceRecord` or `SinkRecord` — by being generic over `ConnectRecord<R>`, the same SMT works in both source and sink pipelines. ### Methods - **`configure(Map<String,?> props)`** — called once at startup. Parse props, usually by constructing a `SimpleConfig`/`AbstractConfig` over your `ConfigDef`, then store typed fields. Validate here and fail fast on bad config. - **`config()`** — return a `ConfigDef` describing each property (name, type, default, importance, doc). Connect uses it for validation and for the REST `validate` endpoint. - **`apply(R record)`** — the hot path. Return a transformed record, the same record unchanged, or **null to drop** it. Must be thread-safe-enough for the task model (one instance per task) and cheap. - **`close()`** — release any resources (rare for stateless SMTs). ## Building the new record `ConnectRecord` is immutable. Use the factory: ``` record.newRecord(record.topic(), record.kafkaPartition(), newKeySchema, newKey, newValueSchema, newValue, record.timestamp(), record.headers()); ``` Return that. To change the topic, pass a different topic string (that's how RegexRouter works). ## The two-schema reality Connect records may be **schemaful** (a `Schema` + a `Struct`) or **schemaless** (`null` schema + a `Map`). A robust SMT branches: - `if (operatingValue == null) return record;` (tombstone safety) - `if (record.valueSchema() == null)` → operate on a `Map` and return a schemaless record. - else → operate on `Struct`/`Schema`; **rebuild the output `Schema` once and cache it** (the JDK pattern is `org.apache.kafka.connect.transforms.util.SchemaUtil` + a `Cache<Schema,Schema>` / `SynchronizedCache`) so you don't rebuild a schema per record — a common performance bug. ## Key vs Value variants The built-ins ship `Key` and `Value` nested subclasses that share an abstract base; the base has an abstract `operatingSchema(record)` / `operatingValue(record)` and `newRecord(...)` that the subclasses implement to point at either the key or the value. Following this convention lets users write `MyTransform$Value`. ## Deployment 1. Build a JAR containing the SMT class (and its deps, isolated). 2. Place it in a directory on the worker's **`plugin.path`** so Connect's classloader isolation picks it up (don't just dump it on the system classpath). 3. Restart the worker; the class is now discoverable. 4. Reference it: `transforms.x.type=com.acme.MyTransform$Value`. ## Gotchas - **Tombstones:** null key/value must not NPE; return the record as-is. - **Schema caching:** rebuilding schemas per record kills throughput. - **Statelessness:** no cross-record state; one record in, one out. - **Immutability:** never mutate the incoming record; always use `newRecord`. - **Config validation:** declare a real `ConfigDef` so misconfig is caught at deploy, not at runtime.
- Why cache the rebuilt schema inside a custom SMT?Schemaful records share a stable input schema; rebuilding the output Schema on every record is expensive. Caching it (e.g. with a SynchronizedCache keyed by input schema) avoids per-record schema construction, a common throughput killer.
- How does a single SMT class work for both source and sink connectors?It is generic over `R extends ConnectRecord<R>`, so the same code handles SourceRecord and SinkRecord. The `record.newRecord(...)` factory produces the correct concrete type.
- Where do you put the JAR so Connect can load a custom SMT?On a directory listed in the worker's `plugin.path`, so Connect's classloader isolation discovers it as a plugin — not on the bare system classpath.
saying these in an interview costs you the question
- Mutating the incoming record instead of returning a new one via record.newRecord — records are immutable.
- Ignoring the schemaless (Map) case and assuming every record is a Struct.
- Rebuilding the output schema per record instead of caching it.
- NPE-ing on tombstones (null value) instead of passing them through.
- Dropping the JAR on the system classpath rather than plugin.path.