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

Axon框架下Kafka消费者未处理事件问题排查求助

Kafka消费者微服务B无响应问题排查

问题概述

微服务A作为Kafka事件生产者,已确认Kafka实例成功存储事件,但微服务B作为消费者仅显示订阅分区,无任何事件处理记录。以下是相关配置与代码:


相关配置与代码

1. Docker Compose配置(Kafka & Zookeeper)

zookeeper:
    container_name: zookeeper-service
    image: 'bitnami/zookeeper:latest'
    ports:
      - '2181:2181'
    environment:
      - ALLOW_ANONYMOUS_LOGIN=yes
    networks:
      - app-tier

kafka:
    container_name: kafka-service
    image: 'bitnami/kafka:latest'
    ports:
      - '9092:9092'
      - '29092:29092'
    volumes:
      - ./kafka-persistence:/bitnami/kafka
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      ALLOW_PLAINTEXT_LISTENER: "yes"
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_LISTENERS: PLAINTEXT://:9092,PLAINTEXT_HOST://0.0.0.0:29092
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
    depends_on:
      - zookeeper
    networks:
      - app-tier

2. 微服务A(生产者)配置与代码

application.yml

axon:
  serializer:
    general: jackson
    messages: jackson
    events: jackson
  kafka:
    bootstrap-servers: kafka:9092
    client-id: myproducer
    default-topic: local.event
    properties:
      security.protocol: PLAINTEXT
    producer:
      transaction-id-prefix: kafka-sample
      retries: 5

Sender类

@Component
public class Sender {

    private static final Logger LOGGER = LoggerFactory.getLogger(Sender.class);

    @Autowired
    private EventBus eventBus;

    public <T> void send(T event) {
        LOGGER.info("publishing event {}", event);
        EventMessage<T> eventMessage = GenericEventMessage.asEventMessage(event);
        eventBus.publish(eventMessage);
    }
}

3. 微服务B(消费者)配置与代码

application.yml

axon:
  serializer:
    general: jackson
    messages: jackson
    events: jackson
  eventhandling:
    processors:
      MyProcessor:
        source: streamableKafkaMessageSource
        mode: TRACKING
        threadCount: 1
        batchSize: 1
  kafka:
    bootstrap-servers: kafka:9092
    client-id: myconsumer
    default-topic: local.event
    properties:
      security.protocol: PLAINTEXT
    consumer:
      event-processor-mode: tracking

KafkaEventConsumer类

@Component
@ProcessingGroup("MyProcessor")
public class KafkaEventConsumer {

    private static final Logger LOGGER = LoggerFactory.getLogger(KafkaEventConsumer.class);

    @EventHandler
    public void handleMyEvent(ObjektfotoDesProjektsFestgelegt myEvent){
        LOGGER.info("got the event {}", myEvent);
        System.out.println("_____ EVENT CONSUMED ______");
    }

}

4. 事件类(A、B微服务均存在)

@Getter
@RequiredArgsConstructor
public class ObjektfotoDesProjektsFestgelegt {
    private final UUID projektId;
    private final String objektfoto;
}

可能的原因及解决方法

1. 事件类型标识不匹配

Axon依赖事件的类型标识路由到对应处理器,若A、B中ObjektfotoDesProjektsFestgelegt类的全限定包路径不一致,或未统一类型别名,消费者会无法识别事件并直接跳过处理。

  • 解决:检查两边事件类的全限定类名完全一致;或给事件类添加@TypeAlias("ObjektfotoDesProjektsFestgelegt")注解,强制统一类型标识。

2. 事务性生产者与消费者隔离级别不匹配

微服务A启用了事务性生产(transaction-id-prefix配置),Kafka事务事件需要消费者配置isolation.level=read_committed才能读取,默认read_uncommitted可能无法获取已提交的事务事件。

  • 解决:在微服务B的配置中添加:
    axon:
      kafka:
        consumer:
          properties:
            isolation.level: read_committed
    
    同时确认生产者的事务已正常提交(Axon在事务上下文发布事件会自动提交,若手动发送需确保事务完成)。

3. 消费者组偏移量问题

若消费者组的偏移量已被提交到最新位置,新启动的消费者会直接从最新偏移量开始,不会消费历史事件;或未显式配置消费者组,导致使用默认组出现异常。

  • 解决:
    1. 在微服务B配置中添加消费者组:
      axon:
        kafka:
          consumer:
            group-id: my-consumer-group
      
    2. 重置消费者组偏移量到最早位置:
      kafka-consumer-groups.sh --bootstrap-server kafka:9092 --group my-consumer-group --topic local.event --reset-offsets --to-earliest --execute
      

4. 事件序列化/反序列化异常

事件类使用@RequiredArgsConstructor但无无参构造函数,Jackson序列化时可能出现异常,消费者默认会静默丢弃事件(无错误日志)。

  • 解决:给事件类添加protected无参构造函数,或配置Axon Jackson序列化支持构造函数注入:
    axon:
      serializer:
        jackson:
          auto-detect-fields: true
          auto-detect-getters-setters: true
          enable-typing: ALWAYS
    

5. 事件处理器配置冲突

微服务B同时配置了eventhandling.processors.MyProcessor.mode=TRACKING和kafka.consumer.event-processor-mode=tracking,虽不冲突,但需确保streamableKafkaMessageSourceBean正确初始化(Axon Kafka Starter默认自动创建,若自定义配置可能遗漏)。

  • 解决:移除重复的模式配置,保留一处即可;检查启动日志确认streamableKafkaMessageSourceBean已成功加载。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 19:56:08