You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

添加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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 05:50:28