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

Kafka Consumer接收异常:JSON无响应、文本报JSON转换错误

Kafka消费JSON消息问题解决

问题现状

使用@InboundChannelAdapter从Kafka Topic轮询消息时:

  • 发送JSON对象{"name" : "foo"}无法触发消费逻辑;
  • 发送普通文本(如sdwadaw)会抛出JSON解析异常。
    需求是将JSON消息转换为CreateResponse类并处理。

问题根源

  1. 参数类型不匹配:@ServiceActivator方法的@Payload参数类型为SftpOutboundFilesDetails,但实际要转换的目标类型是CreateResponse,类型不匹配导致JSON消息无法触发该方法,所以收不到消息。
  2. 非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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 22:25:55