配置ExtractNewRecordState SMT后,如何在Debezium接收Tombstone消息?
Debezium Postgres Connector删除事件丢失问题解答
可以保留Tombstone消息
使用ExtractNewRecordState SMT时,默认会过滤删除事件的Tombstone,但通过添加参数可以保留它。
问题根源
ExtractNewRecordState的默认行为是只提取插入/更新操作的新记录状态,删除事件没有"新状态",所以会被直接丢弃。但Debezium提供了开关来改变这个行为。
解决步骤
给你的unwrap转换添加以下配置项:
"transforms.unwrap.drop.tombstones": "false"
开启这个参数后,SMT会保留删除事件对应的Tombstone消息,不会再过滤掉。
配置里的其他问题
从你给出的配置看,还有几个地方可能出问题:
transforms列表里包含RegexRouter,但没有对应的任何配置参数,这会导致连接器启动失败,直接删掉这个无效项ExtractField$Value要提取的data字段,在经过unwrap处理后可能不存在——因为unwrap已经把原始的Debezium envelope去掉了,值就是记录本身,这时候提取data会得到空值,得确认你的数据结构是否需要这个转换key.converter用的是StringConverter,但key.converter.schemas.enable设为true,StringConverter不支持schema,建议改成false
修正后的完整配置示例
{ "name": "debezium-postgres-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.dbname": "xxx", "database.hostname": "xxx", "database.password": "xx", "database.port": "5432", "database.server.name": "test-server", "database.user": "xxxx", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "plugin.name": "wal2json", "publication.name": "dbz_publication", "slot.name": "debezium_slot", "table.include.list": "myTable", "tasks.max": "1", "transforms": "unwrap,ExtractField,ExtractKey", "transforms.ExtractField.field": "data", "transforms.ExtractField.type": "org.apache.kafka.connect.transforms.ExtractField$Value", "transforms.ExtractKey.field": "key", "transforms.ExtractKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "xxxx", "value.converter.schemas.enable": "true" } }
内容的提问来源于stack exchange,提问作者Sara M.
相关产品推荐
相关产品推荐

