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在事务上下文发布事件会自动提交,若手动发送需确保事务完成)。axon: kafka: consumer: properties: isolation.level: read_committed
3. 消费者组偏移量问题
若消费者组的偏移量已被提交到最新位置,新启动的消费者会直接从最新偏移量开始,不会消费历史事件;或未显式配置消费者组,导致使用默认组出现异常。
- 解决:
- 在微服务B配置中添加消费者组:
axon: kafka: consumer: group-id: my-consumer-group - 重置消费者组偏移量到最早位置:
kafka-consumer-groups.sh --bootstrap-server kafka:9092 --group my-consumer-group --topic local.event --reset-offsets --to-earliest --execute
- 在微服务B配置中添加消费者组:
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
相关产品推荐
相关产品推荐

