添加ExtractField$Key转换后Tombstone消失的原因咨询
Debezium添加ExtractField$Key后墓碑消息消失的问题
问题背景
原本的Debezium MySQL连接器配置能正常发送墓碑消息,但消息key是Struct(id=00000)格式。为了把key改成纯00000的格式,我添加了ExtractField$Key转换,结果key格式正常了,墓碑消息却消失了。
原配置(墓碑消息正常,key为Struct格式)
CREATE SOURCE CONNECTOR `myconn` WITH ( "name" = 'myconn', "connector.class" = 'io.debezium.connector.mysql.MySqlConnector', "tasks.max" = 1, "database.hostname" = 'myconn-db', "database.port" = '${dbPort}', "database.user" = '${dbUsername}', "database.password" = '${dbPassword}', "database.history.kafka.topic" = 'myconn_db_history', "database.history.kafka.bootstrap.servers" = '${bootstrapServer}', "database.server.name" = 'myconn_db', "database.allowPublicKeyRetrieval" = '${allowPublicKeyRetrieval}', "table.include.list" = 'myconn.links,myconn.imports', "message.key.columns" = 'myconn.links:id', "tombstones.on.delete" = true, "null.handling.mode" = 'keep', "transforms" = 'unwrap', "transforms.unwrap.type" = 'io.debezium.transforms.ExtractNewRecordState', "transforms.unwrap.drop.tombstones" = false, "transforms.unwrap.delete.handling.mode" = 'none' );
修改后的配置(key格式正常,墓碑消息消失)
CREATE SOURCE CONNECTOR `myconn` WITH ( "name" = 'myconn', "connector.class" = 'io.debezium.connector.mysql.MySqlConnector', "tasks.max" = 1, "database.hostname" = 'myconn-db', "database.port" = '${dbPort}', "database.user" = '${dbUsername}', "database.password" = '${dbPassword}', "database.history.kafka.topic" = 'myconn_db_history', "database.history.kafka.bootstrap.servers" = '${bootstrapServer}', "database.server.name" = 'myconn_db', "database.allowPublicKeyRetrieval" = '${allowPublicKeyRetrieval}', "table.include.list" = 'myconn.links,myconn.imports', "message.key.columns" = 'myconn.links:id', "tombstones.on.delete" = true, "null.handling.mode" = 'keep', "transforms" = 'unwrap,extractKey', "transforms.unwrap.type" = 'io.debezium.transforms.ExtractNewRecordState', "transforms.unwrap.drop.tombstones" = false, "transforms.unwrap.delete.handling.mode" = 'none', "transforms.extractKey.type" = 'org.apache.kafka.connect.transforms.ExtractField$Key', "transforms.extractKey.field" = 'id', "include.schema.changes" = false );
问题原因
ExtractField$Key转换的默认行为会过滤掉value为null的消息——而Debezium生成的墓碑消息刚好是key为带主键的Struct、value为null的格式。
这个转换默认的null.handling.mode参数值是ignore,一旦检测到消息的value为null,就直接丢弃这条消息,所以墓碑消息就不会被发送到主题了。
解决办法
给extractKey转换添加null.handling.mode配置,设置为pass,让它处理墓碑消息时保留消息,同时正常提取key字段:
修正后的完整配置
CREATE SOURCE CONNECTOR `myconn` WITH ( "name" = 'myconn', "connector.class" = 'io.debezium.connector.mysql.MySqlConnector', "tasks.max" = 1, "database.hostname" = 'myconn-db', "database.port" = '${dbPort}', "database.user" = '${dbUsername}', "database.password" = '${dbPassword}', "database.history.kafka.topic" = 'myconn_db_history', "database.history.kafka.bootstrap.servers" = '${bootstrapServer}', "database.server.name" = 'myconn_db', "database.allowPublicKeyRetrieval" = '${allowPublicKeyRetrieval}', "table.include.list" = 'myconn.links,myconn.imports', "message.key.columns" = 'myconn.links:id', "tombstones.on.delete" = true, "null.handling.mode" = 'keep', "transforms" = 'unwrap,extractKey', "transforms.unwrap.type" = 'io.debezium.transforms.ExtractNewRecordState', "transforms.unwrap.drop.tombstones" = false, "transforms.unwrap.delete.handling.mode" = 'none', "transforms.extractKey.type" = 'org.apache.kafka.connect.transforms.ExtractField$Key', "transforms.extractKey.field" = 'id', "transforms.extractKey.null.handling.mode" = 'pass', -- 新增这行配置 "include.schema.changes" = false );
这样配置后,ExtractField$Key会从墓碑消息的key Struct中提取id作为新key,同时保留value为null的特性,墓碑消息就能正常发送到主题了。
内容的提问来源于stack exchange,提问作者José Vte. Calderón
相关产品推荐
相关产品推荐

