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

无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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 19:54:53