Debezium RabbitMQ Outbox模式下如何实现动态路由键?
问题:Debezium Outbox模式下如何基于MySQL表字段设置RabbitMQ动态路由键
我正尝试通过读取MySQL的Outbox表插入数据,转换为RabbitMQ事件来实现Outbox模式。Outbox表包含以下字段:
- id
- uuid
- payload
- exchange
- routing_key
- produced_timestamp
我使用Debezium及其Outbox事件路由转换器实现该功能,可配置转换器从Outbox表指定事件KEY和负载的列(默认是aggregateid和payload)。目前事件已发送到正确的exchange,但Debezium使用的是debezium.sink.rabbitmq.routingKey指定的静态路由键,而非debezium.transforms.outbox.table.field.event.key配置的表中routing_key字段。查看RabbitMQ消费者源码发现,它仅从上述静态配置项获取路由键,或在设置debezium.sink.rabbitmq.routingKeyFromTopicName时使用topic名称(即exchange)。
请问是否有方法可以基于源表字段,将事件发送到exchange时使用动态路由键?
以下是我的Debezium配置:
# Sink connector config - RabbitMQ debezium.sink.type=rabbitmq debezium.sink.rabbitmq.connection.host=rabbitmq debezium.sink.rabbitmq.connection.port=5672 debezium.sink.rabbitmq.connection.username=guest debezium.sink.rabbitmq.connection.password=guest debezium.sink.rabbitmq.connection.virtual.host=/ debezium.sink.rabbitmq.ackTimeout=3000 debezium.sink.rabbitmq.delivery.mode=2 # Source connector config - MySQL debezium.source.connector.class=io.debezium.connector.mysql.MySqlConnector debezium.source.database.hostname=mysql debezium.source.database.dbname=producer debezium.source.database.port=8779 debezium.source.database.user=root debezium.source.database.password=root debezium.source.table.include.list=producer.outbox debezium.source.database.server.id=184054 debezium.source.offset.storage.file.filename=data/offsets.dat debezium.source.offset.flush.interval.ms=0 debezium.source.topic.prefix=load.test debezium.source.schema.history.internal=io.debezium.storage.file.history.FileSchemaHistory debezium.source.schema.history.internal.file.filename=data/schistory.dat debezium.source.snapshot.mode=initial # Read schema changes only in debezium.source.database.dbname schema.history.internal.store.only.captured.databases.ddl=true # Read and store schema changes on all non-system tables in database schema.history.internal.store.only.captured.tables.ddl=false # Format config debezium.format.key=json debezium.format.value=json # Transformations debezium.transforms=filter,outbox ## FILTER CONFIG # Filter events that are going to be sent to the sink debezium.transforms.filter.type=io.debezium.transforms.Filter # Using Groovy as the langugage for defining the filter debezium.transforms.filter.language=jsr223.groovy # We only want c (create) events (INSERT) debezium.transforms.filter.condition=value.schema().field("op") && value.getString("op") == "c" ## OUTBOX CONFIG # Outbox transformation debezium.transforms.outbox.type=io.debezium.transforms.outbox.EventRouter # Outbox table column to get the destination exchange. Overrides "aggregatetype" debezium.transforms.outbox.route.by.field=exchange # Route to just the value of debezium.transforms.outbox.route.by.field without any prefix debezium.transforms.outbox.route.topic.regex=(?<routedByValue>.*) debezium.transforms.outbox.route.topic.replacement=${routedByValue} # Table unique identifier debezium.transforms.outbox.table.field.event.id=uuid # Event key (routing key for RabbitMQ) debezium.transforms.outbox.table.field.event.key=routing_key # Event timestamp debezium.transforms.outbox.table.field.event.timestamp=produced_timestamp # Event payload debezium.transforms.outbox.table.field.event.payload=payload # Expand json from this: [{"id": 1 to [{"id": 1 debezium.transforms.outbox.table.expand.json.payload=true # Ignore null values during expansion debezium.transforms.outbox.table.json.payload.null.behavior=ignore # Aditional properties to be sent in the event: event_type, produced_timestamp (as timestamp) debezium.transforms.outbox.table.fields.additional.placement=event_type:header,produced_timestamp:header:timestamp # Error if there is no event_type debezium.transforms.outbox.table.fields.additional.error.on.missing=true # Null payload does not emi tombstone event debezium.transforms.outbox.route.tombstone.on.empty.payload=false # If there is an UPDATE to outbox table, logs error and continues consuming debezium.transforms.outbox.table.op.invalid.behavior=error # Quarkus configuration quarkus.log.level=INFO quarkus.log.console.json=false
内容的提问来源于stack exchange,提问作者Victor Castaño Gutierrez
相关产品推荐
相关产品推荐

