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

Spring Cloud Stream中Kafka批量消费CloudEvents异常问题

Spring Cloud Stream Kafka批量读取CloudEvents失败问题解决

问题背景

使用Spring Cloud Stream Kafka绑定器批量读取CloudEvents时,自定义类搭配自定义序列化/反序列化器能正常工作,但改用官方io.cloudevents.kafka组件后,消息无法到达消费者——反序列化看似执行了,但抛出类型转换异常。

相关代码与配置

配置文件

spring:
  cloud:
    function.definition: consumer
    stream:
      bindings:
        producer-out-0:
          destination: audit
          group: audit-producer
          producer:
            useNativeEncoding: true
        consumer-in-0:
          destination: audit
          group: audit-consumer
          consumer:
            batch-mode: true
            useNativeDecoding: true
      kafka:
        binder:
          brokers: localhost:9092
          consumer-properties:
            max.poll.records: 5
            fetch.min.bytes: 10000
            fetch.max.wait.ms: 10000
        bindings:
          producer-out-0:
            producer:
              configuration:
                cloudevents:
                  serializer:
                    encoding: STRUCTURED
                    event_format: application/cloudevents+json
                key.serializer: org.apache.kafka.common.serialization.StringSerializer
#                value.serializer: com.sagar.audit.watcher.domain.MessageSerializer
                value.serializer: io.cloudevents.kafka.CloudEventSerializer
          consumer-in-0:
            consumer:
              configuration:
                key.deserializer: org.apache.kafka.common.serialization.StringDeserializer
#                value.deserializer: com.sagar.audit.watcher.domain.MessageDeserializer
                value.deserializer: io.cloudevents.kafka.CloudEventDeserializer

消费者代码

@Bean
public Consumer<List<CloudEvent>> consumer() {
  System.out.println("inside consumer");
  return s -> s.forEach(auditMessage -> System.out.println("data at loop--" + thread + " -- " + auditMessage));
}

异常信息

2022-10-26 20:31:24.070  WARN [,8289fada18f22581,831ea94d13ef311e] 64368 --- [container-0-C-1] s.c.f.c.c.SmartCompositeMessageConverter : Failure during type conversion by org.springframework.cloud.stream.converter.ApplicationJsonMessageMarshallingConverter@3bf97caf. Will try the next converter.

org.springframework.messaging.converter.MessageConversionException: Could not read JSON: Cannot construct instance of `io.cloudevents.CloudEvent` (no Creators, like default constructor, exist): abstract types either need to be mapped to concrete types, have custom deserializer, or contain additional type information
 at [Source: (String)"[CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-1"}}, extensions={}}, CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-2"}}, extensions={}}]"; line: 1, column: 1]; nested exception is com.fasterxml.jackson.databind.exc.InvalidDefinitionException: Cannot construct instance of `io.cloudevents.CloudEvent` (no Creators, like default constructor, exist): abstract types either need to be mapped to concrete types, have custom deserializer, or contain additional type information
 at [Source: (String)"[CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-1"}}, extensions={}}, CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-2"}}, extensions={}}]"; line: 1, column: 1]
    at org.springframework.messaging.converter.MappingJackson2MessageConverter.convertFromInternal(MappingJackson2MessageConverter.java:237) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.cloud.stream.converter.ApplicationJsonMessageMarshallingConverter.convertFromInternal(ApplicationJsonMessageMarshallingConverter.java:115) ~[spring-cloud-stream-3.2.5.jar:3.2.5]
    at org.springframework.messaging.converter.AbstractMessageConverter.fromMessage(AbstractMessageConverter.java:185) ~[spring-messaging-5.3.23.jar:5.3.23]
    at org.springframework.messaging.converter.AbstractMessageConverter.fromMessage(AbstractMessageConverter.java:176) ~[spring-messaging-5.3.23.jar:5.3.23]
Caused by: com.fasterxml.jackson.databind.exc.InvalidDefinitionException: Cannot construct instance of `io.cloudevents.CloudEvent` (no Creators, like default constructor, exist): abstract types either need to be mapped to concrete types, have custom deserializer, or contain additional type information
 at [Source: (String)"[CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-1"}}, extensions={}}, CloudEvent{id='hello', source=http://localhost, type='example.kafka', datacontenttype='application/json', data=JsonCloudEventData{node={"id":null,"name":"sagar-cloud-2"}}, extensions={}}]"; line: 1, column: 1]
    at com.fasterxml.jackson.databind.exc.InvalidDefinitionException.from(InvalidDefinitionException.java:67) ~[jackson-databind-2.13.4.2.jar:2.13.4.2]

问题原因

异常核心是io.cloudevents.CloudEvent是抽象类,Jackson无法直接实例化它。虽然已经配置了CloudEventDeserializer做Kafka层面的反序列化,但Spring Cloud Stream的消息转换器会对Kafka返回的结果再执行一次JSON转换,导致类型解析失败。

解决方法

方法1:使用具体的CloudEvent实现类

将消费者的泛型参数从抽象类CloudEvent改为SDK提供的具体实现类,比如GenericCloudEvent或CloudEventImpl(取决于你使用的CloudEvents版本):

@Bean
public Consumer<List<GenericCloudEvent>> consumer() {
  System.out.println("inside consumer");
  return s -> s.forEach(auditMessage -> System.out.println("data at loop--" + thread + " -- " + auditMessage));
}

方法2:配置Jackson支持抽象类解析

如果要保留抽象类参数,给Jackson配置CloudEvents专用模块,让它能识别并映射到具体实现类:

@Bean
public ObjectMapper cloudEventObjectMapper() {
    ObjectMapper mapper = new ObjectMapper();
    mapper.registerModule(new CloudEventsJacksonModule());
    return mapper;
}

方法3:禁用多余的消息转换

确认useNativeDecoding: true配置生效,该配置会让Spring Cloud Stream直接使用Kafka的反序列化结果,跳过自身的消息转换。如果仍有问题,可过滤掉不必要的JSON转换器:

@Bean
public SmartCompositeMessageConverter compositeMessageConverter(List<MessageConverter> converters) {
    List<MessageConverter> filteredConverters = converters.stream()
            .filter(converter -> !(converter instanceof ApplicationJsonMessageMarshallingConverter))
            .collect(Collectors.toList());
    return new SmartCompositeMessageConverter(filteredConverters);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 14:01:05