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

Axon-Mongo扩展:能否采用「单租户单集合」方案?

解决方案:Axon Framework + MongoDB 多租户按租户隔离事件集合

核心思路

Axon默认的Mongo事件存储使用固定集合名,要实现按租户分集合,需要自定义MongoTemplate来动态生成集合名,并通过Axon的拦截器传递租户上下文,最后覆盖自动配置的MongoEventStorageEngine。

步骤1:实现租户上下文管理

用ThreadLocal存储当前请求/命令处理线程的租户ID,确保能动态获取:

public class TenantContext {
    private static final ThreadLocal<String> CURRENT_TENANT = new ThreadLocal<>();

    public static void setTenantId(String tenantId) {
        CURRENT_TENANT.set(tenantId);
    }

    public static String getTenantId() {
        return CURRENT_TENANT.get();
    }

    public static void clear() {
        CURRENT_TENANT.remove();
    }
}

步骤2:自定义租户感知的MongoTemplate

重写集合获取逻辑,根据当前租户ID动态生成集合名:

import com.mongodb.client.MongoCollection;
import com.mongodb.client.MongoDatabase;
import org.bson.Document;
import org.axonframework.mongo.MongoTemplate;

public class TenantAwareMongoTemplate implements MongoTemplate {

    private final MongoTemplate delegate;

    public TenantAwareMongoTemplate(MongoTemplate delegate) {
        this.delegate = delegate;
    }

    @Override
    public MongoCollection<Document> domainEventCollection() {
        String tenantId = TenantContext.getTenantId();
        // 按租户ID拼接集合名
        String collectionName = String.format("TENANT-%s-domainevents", tenantId);
        return delegate.database().getCollection(collectionName);
    }

    @Override
    public MongoCollection<Document> snapshotEventCollection() {
        String tenantId = TenantContext.getTenantId();
        String collectionName = String.format("TENANT-%s-snapshotevents", tenantId);
        return delegate.database().getCollection(collectionName);
    }

    // 其他方法直接委托给默认实现
    @Override
    public MongoDatabase database() {
        return delegate.database();
    }

    @Override
    public MongoCollection<Document> trackingTokenCollection() {
        // 跟踪令牌集合无需按租户隔离,复用默认集合
        return delegate.trackingTokenCollection();
    }
}

步骤3:覆盖Axon自动配置,注入自定义组件

创建配置类,替换默认的MongoEventStorageEngine和MongoTemplate,同时添加命令拦截器传递租户上下文:

import org.axonframework.config.EventProcessingConfigurer;
import org.axonframework.mongo.DefaultMongoTemplate;
import org.axonframework.mongo.MongoTemplate;
import org.axonframework.mongo.eventsourcing.eventstore.MongoEventStorageEngine;
import org.axonframework.serialization.Serializer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import com.mongodb.client.MongoClient;

@Configuration
public class AxonMultiTenantConfig {

    @Bean
    public MongoTemplate tenantAwareMongoTemplate(MongoClient mongoClient) {
        // 创建默认MongoTemplate作为委托
        DefaultMongoTemplate defaultTemplate = DefaultMongoTemplate.builder()
                .mongoDatabase(mongoClient)
                .trackingTokenCollection("trackingtokens")
                .build();
        return new TenantAwareMongoTemplate(defaultTemplate);
    }

    @Bean
    public MongoEventStorageEngine mongoEventStorageEngine(MongoTemplate tenantAwareMongoTemplate,
                                                           Serializer defaultSerializer,
                                                           Serializer eventSerializer) {
        // 使用自定义MongoTemplate构建事件存储引擎
        return MongoEventStorageEngine.builder()
                .mongoTemplate(tenantAwareMongoTemplate)
                .eventSerializer(eventSerializer)
                .snapshotSerializer(defaultSerializer)
                .build();
    }

    // 添加命令拦截器,从命令元数据提取租户ID并设置上下文
    @Bean
    public EventProcessingConfigurer configureTenantInterceptor(EventProcessingConfigurer configurer) {
        configurer.registerCommandDispatchInterceptor((config, interceptorChain) -> message -> {
            // 从命令元数据中获取租户ID(需确保命令发送时携带该元数据)
            String tenantId = message.getMetaData().get("tenantId", String.class);
            if (tenantId != null) {
                TenantContext.setTenantId(tenantId);
            }
            try {
                return interceptorChain.proceed(message);
            } finally {
                // 处理完成后清理上下文,避免内存泄漏
                TenantContext.clear();
            }
        });
        return configurer;
    }
}

步骤4:发送命令时携带租户元数据

确保所有命令发送时都附加租户ID元数据,这样拦截器能正确获取上下文:

import org.axonframework.commandhandling.gateway.CommandGateway;
import org.axonframework.messaging.MetaData;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

@Component
public class CommandSender {

    private final CommandGateway commandGateway;

    @Autowired
    public CommandSender(CommandGateway commandGateway) {
        this.commandGateway = commandGateway;
    }

    public void sendYourCommand(Object command, String tenantId) {
        // 附加租户ID到命令元数据
        commandGateway.send(command, MetaData.with("tenantId", tenantId));
    }
}

注意事项

  • MongoDB会自动创建不存在的集合,无需提前为每个租户预创建集合。
  • 若需要隔离Saga数据,可参考上述逻辑自定义MongoSagaStore,替换默认实现。
  • 事件处理线程中若需要租户ID,可通过事件元数据传递(存储事件时将租户ID存入元数据),并添加事件拦截器设置上下文。
  • 务必在拦截器的finally块清理ThreadLocal,防止内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:44:55