如何使用C#反序列化Debezium生成的Kafka消息
Debezium PostgreSQL 变更消息反序列化标准方案
方案一:手动反序列化(无需 Schema Registry)
Debezium 输出的 JSON 消息结构固定,包含 schema 和 payload 两层,其中 payload 的 op 字段标识操作类型(c=创建、u=更新、d=删除、r=快照)。你可以基于这个结构定义对应实体类,用 Jackson 手动完成反序列化。
步骤1:定义实体类
// Debezium 消息包装类 public class DebeziumMessage<T> { private Object schema; // 无需解析 schema 时用 Object 接收即可 private Payload<T> payload; // getter、setter、无参构造器 public static class Payload<T> { private T before; private T after; private String op; // 按需添加 source、ts_ms 等其他字段 // getter、setter、无参构造器 } } // 你已有的 Member 类 public class Member { private Long id; private String name; private Integer age; // getter、setter、toString }
步骤2:消费者配置与消息处理
将 Kafka 消费者的 value.deserializer 设为字符串反序列化器,消费时手动转成实体类:
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class MemberConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put("bootstrap.servers", "kafka-broker:9092"); props.put("group.id", "member-consumer-group"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("your-postgres-topic")); ObjectMapper objectMapper = new ObjectMapper(); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); records.forEach(record -> { try { DebeziumMessage<Member> msg = objectMapper.readValue( record.value(), new TypeReference<DebeziumMessage<Member>>() {} ); String op = msg.getPayload().getOp(); Member member = null; switch (op) { case "c": case "u": member = msg.getPayload().getAfter(); System.out.println("创建/更新会员:" + member); break; case "d": member = msg.getPayload().getBefore(); System.out.println("删除会员:" + member); break; case "r": member = msg.getPayload().getAfter(); System.out.println("快照同步会员:" + member); break; } } catch (Exception e) { // 捕获反序列化异常,避免程序崩溃 System.err.println("消息解析失败:" + record.value()); e.printStackTrace(); } }); } } }
注意:创建操作的 before 为 null,删除操作的 after 为 null,处理时需避免空指针。
方案二:使用 Schema Registry 自动反序列化(生产环境推荐)
Schema Registry 能自动管理消息 schema,避免手动维护结构,适合大规模生产环境。
步骤1:配置 Debezium 连接器
修改连接器配置,将消息序列化为 Avro 格式并注册到 Schema Registry:
name=postgres-connector connector.class=io.debezium.connector.postgresql.PostgresConnector tasks.max=1 database.hostname=postgres-host database.port=5432 database.user=postgres database.password=xxx database.dbname=your-db database.server.name=postgres-server table.include.list=public.member # 关键配置:使用 Avro 转换器 value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://your-schema-registry:8081 key.converter=io.confluent.connect.avro.AvroConverter key.converter.schema.registry.url=http://your-schema-registry:8081
步骤2:消费者配置与代码实现
添加 Maven 依赖:
<dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.4.0</version> <!-- 需与 Confluent 版本匹配 --> </dependency>
生成 Avro 实体类:从 Schema Registry 下载对应主题的 schema,用 Avro 工具生成 Java 类(或用 Maven 插件自动生成):
# 示例命令:下载 schema 并生成类 java -jar avro-tools-1.11.0.jar compile schema schema.avsc ./src/main/java
消费者代码:
import io.confluent.kafka.serializers.KafkaAvroDeserializer; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class AvroMemberConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put("bootstrap.servers", "kafka-broker:9092"); props.put("group.id", "member-avro-consumer-group"); props.put("key.deserializer", KafkaAvroDeserializer.class.getName()); props.put("value.deserializer", KafkaAvroDeserializer.class.getName()); props.put("schema.registry.url", "http://your-schema-registry:8081"); props.put("specific.avro.reader", "true"); // 启用特定类反序列化 // MemberEnvelope 是 Avro 生成的包装类 KafkaConsumer<String, MemberEnvelope> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("your-postgres-topic")); while (true) { ConsumerRecords<String, MemberEnvelope> records = consumer.poll(Duration.ofMillis(100)); records.forEach(record -> { MemberEnvelope envelope = record.value(); String op = envelope.getPayload().getOp(); Member member = null; switch (op) { case "c": case "u": member = envelope.getPayload().getAfter(); break; case "d": member = envelope.getPayload().getBefore(); break; } // 业务逻辑处理 }); } } }
Confluent Json 示例崩溃问题排查
退出码 0x0 通常是正常退出,但如果是意外终止,大概率是以下原因:
- 版本不兼容:Kafka 客户端版本与 Confluent 组件版本不匹配(如 Confluent 7.4 对应 Kafka 3.4),需统一版本。
- 消息格式不匹配:Debezium 用默认 JSON 转换器时,消费者不能用依赖 Schema Registry 的
JsonSchemaDeserializer,需改为字符串反序列化后手动解析,或用无 Schema Registry 的JsonDeserializer。 - 配置错误:Schema Registry 地址配置错误导致连接失败,检查网络可达性与地址正确性。
- 未捕获异常:示例代码未处理反序列化异常,导致异常抛出后程序终止,需添加 try-catch 块捕获
SerializationException。
内容的提问来源于stack exchange,提问作者codingjoe
相关产品推荐
相关产品推荐

