SMTs run in the order listed in transforms, after a source connector or before a sink. Kafka 4.3 129 ships sixteen, among them ValueToKey, ExtractField, ReplaceField, MaskField, Cast, TimestampConverter, RegexRouter and Filter; predicates (TopicNameMatches, HasHeaderKey, RecordIsTombstone) make one conditional. The catalog records have no key, useless for compaction. Three transforms copy id into the key, unwrap it to a number, and reroute the records to the compacted booknest.catalog (Scripting Topic Management):
# Key the catalog by book id and route it to the compacted booknest.catalog topic
C=localhost:33083/connectors
jq '. + {"transforms": "key,id,route",
"transforms.key.type": "org.apache.kafka.connect.transforms.ValueToKey",
"transforms.key.fields": "id",
"transforms.id.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.id.field": "id",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "booknest\\.pg\\.(.*)",
"transforms.route.replacement": "booknest.$1"}' listings/catalog-source.json |
curl -s -X PUT -H 'Content-Type: application/json' -d @- $C/booknest-catalog/config >/dev/null
curl -s -X DELETE $C/booknest-catalog-source # retire the unkeyed copy
sleep 12
docker exec l3-kafka /opt/kafka/bin/kafka-console-consumer.sh --bootstrap-server l3-kafka:9092 \
--topic booknest.catalog --from-beginning --timeout-ms 5000 \
--formatter-property print.key=true 2>/dev/null | cut -c1-72Output
1 {"id":1,"title":"The Quiet Harbor","author":"Mara Ellison","genre":"Fi
5 {"id":5,"title":"The Clockmaker's Paradox","author":"Elena Sokolova","
...A new connector name means fresh offsets, so the current table was copied, book 2 already at 34.50. Keep SMTs stateless; joins and lookups belong in a stream processor.