Axon框架:MySQL存EventStore、MongoDB存Projection重启数据重复问题
问题:Axon框架中MySQL作为EventStore重启后事件重播导致MongoDB数据重复
背景
现有一个将集合存储在MongoDB中的应用,根据文档及视频建议,不应使用MongoDB作为EventStore,转而采用MySQL这类普通RDBMS或Axon Server。已配置使用MySQL作为EventStore、MongoDB在ProductDB中存储Product,但重启SpringBoot服务器后,所有事件会被重播,导致MongoDB中出现重复的Product记录。
配置代码
Application Yaml
spring: data: mongodb: database: productsDB port: 27017 host: localhost mvc: pathmatch: matching-strategy: ant_path_matcher datasource: url: jdbc:mysql://localhost:3306/mysql?autoReconnect=true&useSSL=false&allowPublicKeyRetrieval=true username: dinesh password: abc123 name: mysql jpa: show-sql: true hibernate: ddl-auto: create properties: hibernate: dialect: org.hibernate.dialect.MySQL8Dialect ##Axon configuration axon: serializer: events: jackson general: jackson messages: jackson
Pom.xml(尝试过包含和排除axon-server-connector)
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-mongodb</artifactId> <version>2.7.6</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> <version>2.7.6</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-rest</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <version>${spring.version}</version> <scope>test</scope> </dependency> <dependency> <groupId>io.projectreactor</groupId> <artifactId>reactor-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>io.springfox</groupId> <artifactId>springfox-boot-starter</artifactId> <version>${springfox.version}</version> </dependency> <dependency> <groupId>io.springfox</groupId> <artifactId>springfox-swagger-ui</artifactId> <version>${springfox.version}</version> </dependency> <dependency> <groupId>org.axonframework</groupId> <artifactId>axon-spring-boot-starter</artifactId> <version>${axon-spring-boot-starter.version}</version> <exclusions> <exclusion> <groupId>org.axonframework</groupId> <artifactId>axon-server-connector</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>commons-beanutils</groupId> <artifactId>commons-beanutils</artifactId> <version>${commons-beanutils.version}</version> </dependency> <dependency> <groupId>org.axonframework.extensions.mongo</groupId> <artifactId>axon-mongo</artifactId> <version>${axon-mongo.version}</version> </dependency> <dependency> <groupId>org.axonframework</groupId> <artifactId>axon-metrics</artifactId> <version>${axon-metrics.version}</version> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-jpa</artifactId> </dependency> <dependency> <groupId>com.mysql</groupId> <artifactId>mysql-connector-j</artifactId> <version>8.0.31</version> </dependency>
AxonConfig配置类
@Configuration public class AxonConfig { @Bean public EntityManagerProvider getEntityManagerProvider(){ return new ContainerManagedEntityManagerProvider(); } /* @Bean public EventStore eventStore(EventStorageEngine storageEngine, GlobalMetricRegistry metricRegistry) { return EmbeddedEventStore.builder() .storageEngine(storageEngine) .messageMonitor(metricRegistry .registerEventBus("eventStore")) .build(); }*/ @Bean public TokenStore tokenStore( EntityManagerProvider entityManagerProvider) { return JpaTokenStore.builder().entityManagerProvider(entityManagerProvider) .serializer(JacksonSerializer.defaultSerializer()).build(); } @Bean public EventStorageEngine eventStorageEngine(Serializer serializer, PersistenceExceptionResolver persistenceExceptionResolver, @Qualifier("eventSerializer") Serializer eventSerializer, EntityManagerProvider entityManagerProvider, TransactionManager transactionManager) throws SQLException { JpaEventStorageEngine eventStorageEngine = JpaEventStorageEngine.builder() .snapshotSerializer(serializer) .persistenceExceptionResolver(persistenceExceptionResolver) .eventSerializer(serializer) .entityManagerProvider(entityManagerProvider) .transactionManager(transactionManager) .build(); return eventStorageEngine; } @Bean public EventUpcasterChain eventUpcasters(){ return new EventUpcasterChain(); } }
操作步骤及问题
- POST请求创建Product(正常)
- Axon Server控制台可见事件(正常)
- MongoDB中Product集合及记录正常(正常)
- 停止并重启SpringBoot应用
- MongoDB中出现重复Product记录,因事件被重播
已尝试的方案
- 将EventStore也改为MongoDB(可正常运行,但不符合需求,希望坚持用MySQL/PostgreSQL)
- 尝试添加或移除axon-server-connector(两种情况都未创建DomainEventEntry)
- 启用排除axon-server-connector的配置时,出现如下错误:
2022-12-03 03:24:47.844 WARN 43476 --- [cessor[event]-0] o.a.e.TrackingEventProcessor : Fetch Segments for Processor 'event' failed: org.hibernate.hql.internal.ast.QuerySyntaxException: DomainEventEntry is not mapped [SELECT MIN(e.globalIndex) - 1 FROM DomainEventEntry e]. Preparing for retry in 1s java.lang.IllegalArgumentException: org.hibernate.hql.internal.ast.QuerySyntaxException: DomainEventEntry is not mapped [SELECT MIN(e.globalIndex) - 1 FROM DomainEventEntry e] at org.hibernate.internal.ExceptionConverterImpl.convert(ExceptionConverterImpl.java:138) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.ExceptionConverterImpl.convert(ExceptionConverterImpl.java:181) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.ExceptionConverterImpl.convert(ExceptionConverterImpl.java:188) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.AbstractSharedSessionContract.createQuery(AbstractSharedSessionContract.java:757) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.AbstractSharedSessionContract.createQuery(AbstractSharedSessionContract.java:848) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.AbstractSharedSessionContract.createQuery(AbstractSharedSessionContract.java:114) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[na:na] at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) ~[na:na] at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[na:na] at java.base/java.lang.reflect.Method.invoke(Method.java:566) ~[na:na] at org.springframework.orm.jpa.SharedEntityManagerCreator$SharedEntityManagerInvocationHandler.invoke(SharedEntityManagerCreator.java:311) ~[spring-orm-5.3.24.jar:5.3.24] at com.sun.proxy.$Proxy147.createQuery(Unknown Source) ~[na:na] at org.axonframework.eventsourcing.eventstore.jpa.JpaEventStorageEngine.createTailToken(JpaEventStorageEngine.java:361) ~[axon-eventsourcing-4.6.2.jar:4.6.2] at org.axonframework.eventsourcing.eventstore.AbstractEventStore.createTailToken(AbstractEventStore.java:171) ~[axon-eventsourcing-4.6.2.jar:4.6.2] at org.axonframework.eventhandling.TrackingEventProcessor$WorkerLauncher.lambda$run$1(TrackingEventProcessor.java:1218) ~[axon-messaging-4.6.2.jar:4.6.2] at org.axonframework.common.transaction.TransactionManager.executeInTransaction(TransactionManager.java:47) ~[axon-messaging-4.6.2.jar:4.6.2] at org.axonframework.eventhandling.TrackingEventProcessor$WorkerLauncher.run(TrackingEventProcessor.java:1216) ~[axon-messaging-4.6.2.jar:4.6.2] at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na] Caused by: org.hibernate.hql.internal.ast.QuerySyntaxException: DomainEventEntry is not mapped [SELECT MIN(e.globalIndex) - 1 FROM DomainEventEntry e] at org.hibernate.hql.internal.ast.QuerySyntaxException.generateQueryException(QuerySyntaxException.java:79) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.QueryException.wrapWithQueryString(QueryException.java:103) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.QueryTranslatorImpl.doCompile(QueryTranslatorImpl.java:220) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.QueryTranslatorImpl.compile(QueryTranslatorImpl.java:144) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.engine.query.spi.HQLQueryPlan.<init>(HQLQueryPlan.java:113) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.engine.query.spi.HQLQueryPlan.<init>(HQLQueryPlan.java:73) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.engine.query.spi.QueryPlanCache.getHQLQueryPlan(QueryPlanCache.java:162) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.AbstractSharedSessionContract.getQueryPlan(AbstractSharedSessionContract.java:636) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.internal.AbstractSharedSessionContract.createQuery(AbstractSharedSessionContract.java:748) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] ... 14 common frames omitted Caused by: org.hibernate.hql.internal.ast.QuerySyntaxException: DomainEventEntry is not mapped at org.hibernate.hql.internal.ast.util.SessionFactoryHelper.requireClassPersister(SessionFactoryHelper.java:170) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.tree.FromElementFactory.addFromElement(FromElementFactory.java:91) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.tree.FromClause.addFromElement(FromClause.java:77) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.HqlSqlWalker.createFromElement(HqlSqlWalker.java:334) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.fromElement(HqlSqlBaseWalker.java:3782) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.fromElementList(HqlSqlBaseWalker.java:3671) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.fromClause(HqlSqlBaseWalker.java:746) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.query(HqlSqlBaseWalker.java:602) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.selectStatement(HqlSqlBaseWalker.java:339) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.antlr.HqlSqlBaseWalker.statement(HqlSqlBaseWalker.java:287) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.QueryTranslatorImpl.analyze(QueryTranslatorImpl.java:276) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] at org.hibernate.hql.internal.ast.QueryTranslatorImpl.doCompile(QueryTranslatorImpl.java:192) ~[hibernate-core-5.6.14.Final.jar:5.6.14.Final] ... 20 common frames omitted
解决方案建议
1. 修正JPA实体扫描,生成DomainEventEntry表
出现DomainEventEntry is not mapped错误,是因为Hibernate未识别Axon的JPA实体。可通过两种方式解决:
- 在application.yaml中添加实体扫描路径:
spring: jpa: hibernate: ddl-auto: update # 改为update,避免每次启动重建表 packages-to-scan: - org.axonframework.eventsourcing.eventstore.jpa - 你的业务实体包路径
- 或者在AxonConfig类上添加
@EntityScan注解:
@Configuration @EntityScan(basePackages = {"org.axonframework.eventsourcing.eventstore.jpa", "com.yourpackage"}) public class AxonConfig { // ... 现有代码 }
这样Hibernate会自动创建DomainEventEntry、SnapshotEventEntry和TokenEntry表,分别存储事件、快照和事件处理器进度数据。
2. 确保TokenStore正常跟踪事件处理进度
已配置的JpaTokenStore会通过TokenEntry表记录每个事件处理器的最后处理位置,只要表正常生成,重启应用时事件处理器会从上次中断的位置继续处理,不会重播所有历史事件。
3. 调整事件处理器初始追踪令牌(可选)
如果希望应用启动时直接从最新事件开始处理,无需重播历史,可在AxonConfig中配置:
@Configuration public class AxonConfig { // ... 现有配置 @Autowired public void configureEventProcessor(Configurer configurer) { configurer.eventProcessing() .registerTrackingEventProcessor("event", config -> TrackingEventProcessorConfiguration.forSingleThreadedProcessing() .initialTrackingToken(streamableMessageSource -> streamableMessageSource.createHeadToken())); } }
注意替换"event"为你的事件处理器实际名称。
4. 实现投影逻辑的幂等性
如果无法完全避免事件重播,需让MongoDB投影逻辑具备幂等性。保存Product时,检查是否已存在相同ID的记录,存在则更新而非插入:
@Service @ProcessingGroup("event") public class ProductProjection { private final MongoTemplate mongoTemplate; public ProductProjection(MongoTemplate mongoTemplate) { this.mongoTemplate = mongoTemplate; } @EventHandler public void handle(ProductCreatedEvent event) { Query query = Query.query(Criteria.where("id").is(event.getProductId())); Update update = Update.update("name", event.getName()) .set("price", event.get
相关产品推荐
相关产品推荐

