Kafka Connect ToLogicalTopicRouter (Debezium) SMT Usage Reference for Confluent Cloud
The ToLogicalTopicRouter SMT (io.debezium.transforms.ToLogicalTopicRouter)
reroutes Debezium change event records from multiple physical tables, such as
per-shard source tables, into a single logical topic.
Description
Use this SMT to combine change events that originate from sharded or partitioned
tables, such as customers_shard1 and customers_shard2, into one topic.
The SMT also adds a field to the record
key that identifies the originating table so that keys remain unique after the
tables merge into a single topic.
ToLogicalTopicRouter is the current Debezium implementation of topic
routing. It supersedes the older ByLogicalTableRouter SMT
(io.debezium.transforms.ByLogicalTableRouter), which is Debezium deprecated.
The two SMTs provide similar configuration options and routing behavior, which
Debezium documents under the former name, but some option names and defaults
differ. For example, the default value of key.field.name is
__dbz__sourceIdentifier for ToLogicalTopicRouter and
__dbz__physicalTableIdentifier for ByLogicalTableRouter, and the cache
size option is named logical.destination.cache.size instead of
logical.table.cache.size. For complete details, see the official
Debezium Topic Routing SMT
documentation. For other topic-routing SMTs, see
Kafka Connect Single Message Transformation Reference for Confluent Cloud.
Note
When using this SMT, consider the following:
Keep
key.enforce.uniquenessset totrue(the default) so the SMT adds a source identifier to the record key, keeping keys unique after merging multiple tables into one topic.Ensure the merged tables have compatible value schemas. Routing tables with incompatible schemas can cause schema compatibility errors in Schema Registry.
Set
key.field.nameto a name that doesn’t already exist in the record key to avoid a field collision.
Limitations
The ToLogicalTopicRouter SMT is available only for managed Debezium change
data capture (CDC) Source connectors, such as PostgreSQL, MySQL, SQL Server, and
MariaDB, and the Oracle XStream CDC Source connector.
Properties
The following table lists the configuration options for the
ToLogicalTopicRouter SMT.
Name |
Description |
Type |
Default |
Importance |
|---|---|---|---|---|
|
A regular expression that is applied to the topic name of each change event record to determine whether the record should be rerouted to a logical topic. This option is required. |
string |
low |
|
|
A replacement string that specifies the logical destination topic
name. It can reference groups captured by |
string |
low |
|
|
Whether to add a field to the record’s change event key that identifies the originating physical source. This keeps keys unique when multiple tables are merged into a single topic. |
boolean |
|
low |
|
The name of the field that is added to the change event key to identify the originating source. |
string |
|
low |
|
A regular expression that is applied to the record’s original topic name to determine the value that is inserted into the key field. |
string |
low |
|
|
A replacement string that determines the value of the inserted key
field, using groups captured by |
string |
low |
|
|
Specifies how to adjust the schema name of the message key for
converter compatibility. Accepts |
string |
|
low |
|
The maximum number of entries kept in the least recently used (LRU) cache that the SMT uses for schema and regular-expression caching. |
int |
|
low |
Example
The following example merges change events from tables that match
(.*)customers_shard[0-9]+ into a single $1customers_all topic and adds
the source identifier to the record key.
"transforms": "route",
"transforms.route.type": "io.debezium.transforms.ToLogicalTopicRouter",
"transforms.route.topic.regex": "(.*)customers_shard[0-9]+",
"transforms.route.topic.replacement": "$1customers_all",
"transforms.route.key.field.name": "__dbz__sourceIdentifier",
"transforms.route.schema.name.adjustment.mode": "avro"
Predicates
Transformations can be configured with predicates so that the transformation is applied only to records which satisfy a condition. You can use predicates in a transformation chain and, when combined with the Kafka Connect Filter (Kafka) SMT Usage Reference for Confluent Cloud, predicates can conditionally filter out specific records. For details and examples, see Predicates.