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

Confluent Cloud中Debezium与KSQLDB的墓碑消息及Struct键处理咨询

处理Debezium墓碑消息与KSQLDB表匹配的问题

背景

我通过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个转换,这种方式存在局限性。

问题

  1. 能否通过配置Debezium(或任意Kafka Connect连接器)直接发送0000而非Struct(id=0000)格式的键,无需使用转换?
  2. 处理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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 06:15:37