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.uniqueness set to true (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.name to 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

topic.regex

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

topic.replacement

A replacement string that specifies the logical destination topic name. It can reference groups captured by topic.regex. This option is required.

string

low

key.enforce.uniqueness

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

true

low

key.field.name

The name of the field that is added to the change event key to identify the originating source.

string

__dbz__sourceIdentifier

low

key.field.regex

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

key.field.replacement

A replacement string that determines the value of the inserted key field, using groups captured by key.field.regex.

string

low

schema.name.adjustment.mode

Specifies how to adjust the schema name of the message key for converter compatibility. Accepts none, avro, or avro_unicode.

string

none

low

logical.destination.cache.size

The maximum number of entries kept in the least recently used (LRU) cache that the SMT uses for schema and regular-expression caching.

int

16

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.