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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 11:28:12