Kafka Consumer接收异常:JSON无响应、文本报JSON转换错误
Kafka消费JSON消息问题解决
问题现状
使用@InboundChannelAdapter从Kafka Topic轮询消息时:
- 发送JSON对象
{"name" : "foo"}无法触发消费逻辑; - 发送普通文本(如
sdwadaw)会抛出JSON解析异常。
需求是将JSON消息转换为CreateResponse类并处理。
问题根源
- 参数类型不匹配:
@ServiceActivator方法的@Payload参数类型为SftpOutboundFilesDetails,但实际要转换的目标类型是CreateResponse,类型不匹配导致JSON消息无法触发该方法,所以收不到消息。 - 非JSON消息解析失败:
StringJsonMessageConverter会尝试将所有消息内容解析为JSON,普通文本不符合JSON语法,因此抛出JsonParseException。
修复方案
1. 修正ServiceActivator参数类型
将消费方法的参数类型改为CreateResponse,确保转换后的对象能匹配方法签名:
@ServiceActivator(inputChannel = "inputChannel") void consumeIt(@Payload CreateResponse cr, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) throws JSchException, URISyntaxException, IOException { log.info("In SERVICE ACTIVATOR "); MDC.put("transaction.id", String.valueOf(UUID.randomUUID())); fileServices.processFile(cr); acknowledgment.acknowledge(); log.info("ACKNOWLEDGED: "); }
2. 优化消息转换配置(二选一)
方案一:使用JsonDeserializer(推荐)
直接通过Kafka的JsonDeserializer完成JSON到CreateResponse的转换,简化配置逻辑:
@Bean public ConsumerFactory<String, CreateResponse> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); // 替换为JsonDeserializer处理值的反序列化 props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 60000); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 10); props.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, "10000"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "30000"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "myGroupID"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); // 指定默认的反序列化类型为CreateResponse props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, CreateResponse.class); // 信任所有包(根据实际情况调整) props.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); return new DefaultKafkaConsumerFactory<>(props); } @Bean @InboundChannelAdapter(channel = "inputChannel", poller = @Poller(fixedDelay = "5000")) // 调整泛型类型为<String, CreateResponse> public KafkaMessageSource<String, CreateResponse> consumeMsg(ConsumerFactory<String, CreateResponse> consumerFactory) { KafkaMessageSource<String, CreateResponse> kafkaMessageSource = new KafkaMessageSource<>(consumerFactory, new ConsumerProperties("myTopic")); kafkaMessageSource.getConsumerProperties().setGroupId("myGroupId"); kafkaMessageSource.getConsumerProperties().setClientId("myClientId"); return kafkaMessageSource; } // 移除原StringJsonMessageConverter的Bean配置
方案二:保留StringJsonMessageConverter并添加异常处理
如果需要保留原字符串反序列化+消息转换器的模式,添加异常处理来忽略非JSON消息:
// 保留原messageConverter和consumerFactory配置(使用StringDeserializer) // 修正ServiceActivator方法,添加异常处理通知 @ServiceActivator(inputChannel = "inputChannel", adviceChain = "jsonErrorHandler") void consumeIt(@Payload CreateResponse cr, @Header(KafkaHeaders.ACKNOWLEDGMENT) Acknowledgment acknowledgment) throws JSchException, URISyntaxException, IOException { log.info("In SERVICE ACTIVATOR "); MDC.put("transaction.id", String.valueOf(UUID.randomUUID())); fileServices.processFile(cr); acknowledgment.acknowledge(); log.info("ACKNOWLEDGED: "); } // 定义异常处理通知,忽略JSON解析异常 @Bean public Advice jsonErrorHandler() { return new MethodInterceptor() { @Override public Object invoke(MethodInvocation invocation) throws Throwable { try { return invocation.proceed(); } catch (ConversionException e) { log.warn("收到非JSON格式消息,跳过处理", e); // 手动确认消息,避免重复消费 Acknowledgment acknowledgment = (Acknowledgment) invocation.getArguments()[1]; acknowledgment.acknowledge(); return null; } } }; }
3. 确保CreateResponse类符合序列化要求
CreateResponse类需要具备无参构造函数,以及与JSON字段对应的属性和getter/setter:
public class CreateResponse { private String name; // 必须有无参构造函数 public CreateResponse() {} public String getName() { return name; } public void setName(String name) { this.name = name; } }
异常说明
发送普通文本时抛出的JsonParseException,是因为StringJsonMessageConverter强制将所有消息解析为JSON格式,而普通文本不符合JSON语法规则。通过上述异常处理方案,可以跳过这类无效消息,不影响正常JSON消息的消费。
内容的提问来源于stack exchange,提问作者A S
相关产品推荐
相关产品推荐

