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

如何规避死信主题的Schema Registry并配置不同反序列化器?

问题1:将Avro消息以JSON格式发送至死信主题(无需Schema Registry)

当前问题是错误处理器默认复用原有KafkaTemplate的Avro序列化器,导致无法将Avro对象直接序列化为JSON。解决核心是为死信发送单独配置JSON序列化的KafkaTemplate,并在死信恢复器中完成Avro到JSON的格式转换。

步骤1:配置JSON序列化专属KafkaTemplate

创建专门用于死信发送的ProducerFactory和KafkaTemplate:

@Bean("dltJsonProducerFactory")
public ProducerFactory<String, Object> dltJsonProducerFactory(KafkaProperties kafkaProperties) {
    Map<String, Object> configs = kafkaProperties.buildProducerProperties();
    // 替换为JSON序列化器
    configs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
    // 关闭类型头(无需依赖Schema Registry)
    configs.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);
    return new DefaultKafkaProducerFactory<>(configs);
}

@Bean("dltJsonKafkaTemplate")
public KafkaTemplate<String, Object> dltJsonKafkaTemplate(
        @Qualifier("dltJsonProducerFactory") ProducerFactory<String, Object> producerFactory) {
    return new KafkaTemplate<>(producerFactory);
}

步骤2:自定义死信恢复器完成Avro转JSON

修改ErrorHandler,使用上述JSON模板,并在发送前将Avro对象转换为JSON格式:

@Bean
public DefaultErrorHandler errorHandler(@Qualifier("dltJsonKafkaTemplate") KafkaTemplate<String, Object> dltJsonTemplate) {
    DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(dltJsonTemplate,
            // 指定死信主题规则(原主题后缀加.dlt)
            (record, ex) -> new TopicPartition(record.topic() + ".dlt", record.partition()),
            // 转换Avro对象为JSON字符串
            (record, ex) -> {
                Object avroValue = record.value();
                try (ByteArrayOutputStream out = new ByteArrayOutputStream()) {
                    // 用Avro原生JsonEncoder完成格式转换
                    JsonEncoder encoder = EncoderFactory.get().jsonEncoder(avroValue.getSchema(), out);
                    DatumWriter<Object> writer = new GenericDatumWriter<>(avroValue.getSchema());
                    writer.write(avroValue, encoder);
                    encoder.flush();
                    String jsonStr = out.toString(StandardCharsets.UTF_8);
                    // 返回转换后的ProducerRecord
                    return new ProducerRecord<>(record.topic() + ".dlt", record.key(), jsonStr);
                } catch (IOException e) {
                    throw new RuntimeException("Avro转JSON失败", e);
                }
            });

    DefaultErrorHandler handler = new DefaultErrorHandler(recoverer, new FixedBackOff(0L, 0L));
    handler.addNotRetryableExceptions(IllegalArgumentException.class, SaveFactFailedException.class);
    return handler;
}
问题2:为不同监听器配置独立反序列化器

全局配置会统一应用到所有监听器,需为每个监听器创建独立的ConsumerFactory和ContainerFactory,在@KafkaListener注解中指定对应的容器工厂即可实现差异化配置。

步骤1:创建Avro反序列化容器工厂(原监听器使用)

@Bean("avroConsumerFactory")
public ConsumerFactory<String, Object> avroConsumerFactory(KafkaProperties kafkaProperties) {
    Map<String, Object> configs = kafkaProperties.buildConsumerProperties();
    // 保留原有ErrorHandlingDeserializer和Avro反序列化配置
    configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    configs.put(ErrorHandlingDeserializer.KEY_DELEGATE_CLASS, StringDeserializer.class.getName());
    configs.put(ErrorHandlingDeserializer.VALUE_DELEGATE_CLASS, AWSKafkaAvroDeserializer.class.getName());
    return new DefaultKafkaConsumerFactory<>(configs);
}

@Bean("avroListenerContainerFactory")
public ConcurrentKafkaListenerContainerFactory<String, Object> avroListenerContainerFactory(
        @Qualifier("avroConsumerFactory") ConsumerFactory<String, Object> consumerFactory) {
    ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    return factory;
}

步骤2:创建JSON反序列化容器工厂(DLT监听器使用)

@Bean("jsonConsumerFactory")
public ConsumerFactory<String, String> jsonConsumerFactory(KafkaProperties kafkaProperties) {
    Map<String, Object> configs = kafkaProperties.buildConsumerProperties();
    // 替换为JSON反序列化器
    configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
    // 设置信任包(根据业务调整,避免安全风险)
    configs.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
    // 关闭类型头(DLT消息无Avro类型信息)
    configs.put(JsonDeserializer.USE_TYPE_INFO_HEADERS, false);
    return new DefaultKafkaConsumerFactory<>(configs);
}

@Bean("jsonListenerContainerFactory")
public ConcurrentKafkaListenerContainerFactory<String, String> jsonListenerContainerFactory(
        @Qualifier("jsonConsumerFactory") ConsumerFactory<String, String> consumerFactory) {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    return factory;
}

步骤3:在监听器中指定对应容器工厂

// 原Avro消息监听器,绑定Avro容器工厂
@KafkaListener(topics = "original-topic", groupId = "mobile-layout-group", containerFactory = "avroListenerContainerFactory")
public void consumeAvroMessage(ConsumerRecord<String, Object> record) {
    // 原有业务处理逻辑
}

// DLT JSON消息监听器,绑定JSON容器工厂
@KafkaListener(topics = "original-topic.dlt", groupId = "dlt-group", containerFactory = "jsonListenerContainerFactory")
public void consumeDltJsonMessage(ConsumerRecord<String, String> record) {
    // 处理JSON格式死信消息,可自行将字符串转换为业务对象
}

内容的提问来源于stack exchange,提问作者user3236794

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 08:50:30