Confluent Cloud中Debezium与KSQLDB的墓碑消息及Struct键处理咨询
背景
我通过KSQLDB创建了Debezium Kafka连接器,当数据库表中某行被删除时,Debezium发送的墓碑消息格式如下:
KEY: Struct(cliente_cod=0000) | BODY: null
而KSQLDB中物化表的列示例为:
ID: 0000 | NAME: xxxx | SURNAME: xxxx
未做转换时,墓碑消息的Struct类型键Struct(cliente_cod=0000)无法与表中的ID值0000匹配,导致对应行无法被删除。若将Struct类型直接存为表ID,又会给表关联操作带来麻烦。
尝试过用PARTITION BY重新分区流,但流不识别墓碑消息(这是物化视图的概念),会忽略null内容,所以此方法无效。
目前可行的方案是在KSQLDB连接器定义中添加转换配置,示例如下:
"transforms.extractClienteKey.type" = 'org.apache.kafka.connect.transforms.ExtractField$Key', "transforms.extractClienteKey.field" = 'cliente_cod', "transforms.extractClienteKey.predicate" = 'IsClienteTopic',
转换后墓碑消息格式变为:
KEY: 0000 | BODY: null
但数据库中有大量主键名称不同的表(比如30个表,主键名包括client_id、user_id等),需要按主题为每个表配置不同的ExtractField$Key转换。而Confluent Cloud限制每个连接器最多配置10个转换,这种方式存在局限性。
问题
- 能否通过配置Debezium(或任意Kafka Connect连接器)直接发送
0000而非Struct(id=0000)格式的键,无需使用转换? - 处理Debezium墓碑消息与KSQLDB表的正确方式是什么?转换是唯一方案吗?是否有其他替代方法?
问题1解答
可以通过Debezium的原生配置直接调整键的格式,无需依赖Kafka Connect转换:
单主键表场景:如果数据库表只有单个主键列,可通过关闭键转换器的Schema启用配置,让Debezium直接发送主键原始值。示例配置:
"key.converter" = "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable" = "false",注意:该配置会作用于连接器同步的所有表,若存在复合主键的表,此方式不适用(复合主键无法简化为单个值)。
表级覆盖配置:Debezium支持按表单独配置,可针对特定单主键表关闭Schema启用,示例:
"table.include.list" = "db.clientes,db.users", "db.clientes.key.converter.schemas.enable" = "false", "db.users.key.converter.schemas.enable" = "false"同样,此方式仅适用于单主键表,复合主键仍会以Struct形式存在。
如果是复合主键场景,目前没有直接配置能将其转换为单个值,仍需依赖ExtractField转换。
问题2解答
转换是常用方案,但并非唯一选择,以下是几种替代方法:
1. KSQLDB中直接处理结构化键
无需修改连接器配置,在KSQLDB创建表时将键定义为STRUCT类型,直接提取字段匹配主键:
CREATE TABLE clientes ( cliente_cod STRING PRIMARY KEY, NAME STRING, SURNAME STRING ) WITH ( KAFKA_TOPIC = 'db.clientes', VALUE_FORMAT = 'AVRO', KEY_FORMAT = 'AVRO' -- 需与Debezium发送的键格式一致 );
此时墓碑消息的Struct键会被正确解析,KSQLDB能识别cliente_cod字段并匹配主键,完成行删除。关联表时需显式提取Struct字段,例如:
SELECT c.cliente_cod, u.user_name FROM clientes c JOIN users u ON c.cliente_cod = u.user_id;
2. 自定义Debezium消息转换器
开发或使用第三方自定义转换器,在消息生成阶段直接将键转换为所需格式。比如基于Debezium的io.debezium.transforms.ExtractNewRecordState扩展,或自定义逻辑批量处理不同主键名称的表,避免配置多个ExtractField转换。此方式适合具备开发能力的团队。
3. 拆分连接器
针对Confluent Cloud的转换数量限制,可将大连接器拆分为多个小连接器,按主键类型或业务分组处理表。例如一个连接器处理所有主键为client_id的表,另一个处理user_id的表,确保每个连接器的转换数量不超过10个。
总结
转换是最直接的方案,但根据场景不同,也可选择在KSQLDB中处理结构化键、自定义转换器或拆分连接器。单主键表优先使用Debezium的key.converter.schemas.enable=false配置;复合主键或多主键名称场景,可结合实际需求选择转换、KSQLDB结构化处理或拆分连接器的方案。
内容的提问来源于stack exchange,提问作者José Vte. Calderón

