如何规避死信主题的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
相关产品推荐
相关产品推荐

