无Spring依赖下,基于Kafka+PostgreSQL配置AxonFramework的问题
Axon无Spring依赖下整合Kafka与PostgreSQL配置方案
背景
已定义AxonConnector接口用于访问Spring环境中Axon发布到Kafka并存储在PostgreSQL的事件,基于AxonServer的实现已正常运行。当前需在无Spring依赖场景下,完成PostgreSQL作为EventStore、Kafka作为消息中间件的Axon整合,已完成PostgreSQL EventStore配置和Kafka生产者/消费者初始化,需后续整合步骤。
首先,确保已引入Axon Kafka扩展依赖(Maven为例):
<dependency> <groupId>org.axonframework.extensions.kafka</groupId> <artifactId>axon-kafka</artifactId> <version>4.9.0</version> <!-- 匹配你的Axon核心版本 --> </dependency>
1. 配置Kafka消息序列化器
Axon需要专门的序列化器处理领域事件/命令,不能直接用StringSerializer,需配置Axon的JacksonSerializer适配Kafka:
// 复用已定义的serializer(JacksonSerializer) Serializer axonSerializer = JacksonSerializer.defaultSerializer(); // Kafka消息转换器,将Axon消息转为Kafka的ProducerRecord KafkaMessageConverter<String, Object> messageConverter = DefaultKafkaMessageConverter.builder() .keySerializer(axonSerializer) .valueSerializer(axonSerializer) .build();
2. 配置Kafka事件发布器(EventPublisher)
将Axon的领域事件发布到Kafka主题:
// 基于已初始化的KafkaProducer创建KafkaPublisher KafkaPublisher<String, Object> kafkaPublisher = KafkaPublisher.<String, Object>builder() .kafkaProducer(producer) .messageConverter(messageConverter) .defaultTopic("axon-events") // 对应Spring端Axon发布事件的Kafka主题 .build(); // 将KafkaPublisher注册为Axon的EventBus的事件处理器,实现事件发布到Kafka configurer.registerComponent(EventProcessor.class, config -> { SimpleEventBus eventBus = (SimpleEventBus) config.eventBus(); eventBus.registerHandlerInterceptor(kafkaPublisher); return kafkaPublisher; });
3. 配置Kafka消息处理器(消费Kafka中的命令/事件)
从Kafka消费消息并路由到Axon的命令总线/事件总线:
// 创建KafkaMessageSource,负责从Kafka消费消息并转为Axon消息 KafkaMessageSource<String, Object> kafkaMessageSource = KafkaMessageSource.<String, Object>builder() .kafkaConsumer(consumer) .messageConverter(messageConverter) .build(); // 将Kafka消费的命令路由到Axon的CommandBus configurer.registerComponent(CommandBus.class, config -> { SimpleCommandBus commandBus = (SimpleCommandBus) config.commandBus(); // 消费Kafka中的命令消息,发送到Axon命令总线 Executors.newSingleThreadExecutor().submit(() -> { while (!Thread.currentThread().isInterrupted()) { kafkaMessageSource.poll(messages -> { messages.forEach(message -> commandBus.dispatch(message)); }); } }); return commandBus; }); // 若需要消费事件,可类似配置将Kafka事件路由到EventBus // configurer.registerComponent(EventBus.class, config -> { // SimpleEventBus eventBus = (SimpleEventBus) config.eventBus(); // Executors.newSingleThreadExecutor().submit(() -> { // while (!Thread.currentThread().isInterrupted()) { // kafkaMessageSource.poll(messages -> { // messages.forEach(eventBus::publish); // }); // } // }); // return eventBus; // });
4. 完善AxonConnector接口实现
将Axon核心组件与Kafka整合逻辑封装到AxonConnector实现类中:
public class KafkaPostgreAxonConnector implements AxonConnector { private final Configuration axonConfig; private final KafkaPublisher<String, Object> kafkaPublisher; private final ExecutorService kafkaConsumerExecutor; // 构造函数注入已配置好的axonConfig、kafkaPublisher、消费线程池等 public KafkaPostgreAxonConnector(Configuration axonConfig, KafkaPublisher<String, Object> kafkaPublisher, ExecutorService kafkaConsumerExecutor) { this.axonConfig = axonConfig; this.kafkaPublisher = kafkaPublisher; this.kafkaConsumerExecutor = kafkaConsumerExecutor; } @Override public void start() { // Axon配置已通过configurer.start()启动,这里可添加Kafka连接验证 System.out.println("Axon-Kafka-PostgreSQL Connector started"); } @Override public void close() throws Exception { kafkaConsumerExecutor.shutdown(); kafkaPublisher.close(); axonConfig.shutdown(); } @Override public void registerEventHandler(Object eventHandler) { axonConfig.eventBus().register(eventHandler); } @Override public <T> void registerAggregate(Class<T> aggregate) { axonConfig.commandBus().subscribe(aggregate.getSimpleName(), axonConfig.aggregateFactory(aggregate)::handle); } @Override public void sendCommand(Object command) { axonConfig.commandBus().dispatch(GenericCommandMessage.asCommandMessage(command)); } @Override public <R> R sendQuery(Object query) { return axonConfig.queryBus().query(GenericQueryMessage.asQueryMessage(query), ResponseTypes.instanceOf(Object.class)).join(); } }
5. 完整初始化流程
整合所有配置,创建AxonConnector实例:
public class AxonConnectorBootstrap { public static void main(String[] args) { // 1. 初始化PostgreSQL EventStore(复用你已有的代码) HikariConfig hikariConfig = new HikariConfig(); hikariConfig.setJdbcUrl("jdbc:postgresql://localhost:5432/axon_db"); hikariConfig.setUsername("postgres"); hikariConfig.setPassword("password"); // ... 其他Hikari配置 DataSource dataSource = new HikariDataSource(hikariConfig); ConnectionProvider connectionProvider = new UnitOfWorkAwareConnectionProviderWrapper( new DataSourceConnectionProvider(dataSource) ); Serializer serializer = JacksonSerializer.defaultSerializer(); EventStorageEngine eventStorageEngine = JdbcEventStorageEngine.builder() .connectionProvider(connectionProvider) .snapshotFilter(SnapshotFilter.allowAll()) .eventSerializer(serializer) .snapshotSerializer(serializer) .build(); EmbeddedEventStore eventStore = EmbeddedEventStore.builder() .storageEngine(eventStorageEngine) .build(); // 2. 初始化Kafka生产者/消费者(复用你已有的代码) Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); KafkaProducer<String, Object> producer = new KafkaProducer<>(producerProps); Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "axon-connector-group"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); KafkaConsumer<String, Object> consumer = new KafkaConsumer<>(consumerProps); consumer.subscribe(Collections.singletonList("axon-commands")); // 3. 配置Axon与Kafka整合 Configurer configurer = DefaultConfigurer.defaultConfiguration() .configureEventStore(c -> eventStore) .configureSerializer(c -> serializer) .configureEventSerializer(c -> serializer) .configureMessageSerializer(c -> serializer); // 配置Kafka消息转换器 KafkaMessageConverter<String, Object> messageConverter = DefaultKafkaMessageConverter.builder() .keySerializer(serializer) .valueSerializer(serializer) .build(); // 配置Kafka事件发布器 KafkaPublisher<String, Object> kafkaPublisher = KafkaPublisher.<String, Object>builder() .kafkaProducer(producer) .messageConverter(messageConverter) .defaultTopic("axon-events") .build(); configurer.registerComponent(EventProcessor.class, config -> { ((SimpleEventBus) config.eventBus()).registerHandlerInterceptor(kafkaPublisher); return kafkaPublisher; }); // 配置Kafka命令消费 KafkaMessageSource<String, Object> kafkaMessageSource = KafkaMessageSource.<String, Object>builder() .kafkaConsumer(consumer) .messageConverter(messageConverter) .build(); ExecutorService consumerExecutor = Executors.newSingleThreadExecutor(); configurer.registerComponent(CommandBus.class, config -> { SimpleCommandBus commandBus = (SimpleCommandBus) config.commandBus(); consumerExecutor.submit(() -> { while (!Thread.currentThread().isInterrupted()) { kafkaMessageSource.poll(messages -> messages.forEach(commandBus::dispatch)); } }); return commandBus; }); // 启动Axon配置 Configuration axonConfig = configurer.start(); // 4. 创建AxonConnector实例 AxonConnector connector = new KafkaPostgreAxonConnector(axonConfig, kafkaPublisher, consumerExecutor); connector.start(); // 示例:注册事件处理器、发送命令 connector.registerEventHandler(new MyEventHandler()); connector.sendCommand(new MyCommand("test-id", "test-data")); } // 示例事件处理器 static class MyEventHandler { @EventHandler public void handle(MyEvent event) { System.out.println("Received event: " + event); } } }
关键注意事项
- 版本匹配:确保Axon核心版本与Axon Kafka扩展版本一致,避免兼容性问题。
- 序列化一致性:Spring端与无Spring端必须使用相同的序列化配置(如Jackson的模块、类型信息),否则消息无法正确反序列化。
- Kafka主题对齐:消费/发布的Kafka主题需与Spring端Axon配置的主题完全一致。
- 资源关闭:在
AutoCloseable的close()方法中需正确关闭Kafka生产者/消费者、Axon配置、线程池等资源,避免内存泄漏。
内容的提问来源于stack exchange,提问作者CronixZero
相关产品推荐
相关产品推荐

