在实际工作中,我遇到了这种连接多个mongodb数据源的问题,真实情况是让我将6个mongodb的20个库的数据同步到mysql,下面我将演示如何实现这一功能(用三个数据源来演示)

1.添加依赖

        <!--mongodb java驱动-->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-mongodb</artifactId>
        </dependency>

2.在springboot工程中自定义MongoTemplate

package com.ct.sync.config;
 
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.mongo.MongoProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.SimpleMongoClientDatabaseFactory;
import org.springframework.data.mongodb.repository.config.EnableMongoRepositories;
 
@Configuration
@EnableMongoRepositories(
        basePackages = "com.ct.sync.domain.*",
        mongoTemplateRef = "oneMongo")
public class OneMongoTemplate {
 
    @Autowired
    @Qualifier("oneMongoProperties")
    private MongoProperties mongoProperties;
 
    @Primary
    @Bean(name = "oneMongo")
    public MongoTemplate oneMongoTemplate() throws Exception {
        return new MongoTemplate(oneFactory(this.mongoProperties));
    }
 
    @Bean
    @Primary
    public SimpleMongoClientDatabaseFactory oneFactory(MongoProperties mongoProperties) throws Exception {
        return new SimpleMongoClientDatabaseFactory(mongoProperties.getUri());
    }
}
 
package com.ct.sync.config;
 
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.mongo.MongoProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.SimpleMongoClientDatabaseFactory;
import org.springframework.data.mongodb.repository.config.EnableMongoRepositories;
 
@Configuration
@EnableMongoRepositories(
        basePackages = "com.ct.sync.domain.*",
        mongoTemplateRef = "twoMongo")
public class TwoMongoTemplate {
 
    @Autowired
    @Qualifier("twoMongoProperties")
    private MongoProperties mongoProperties;
 
    @Bean(name = "twoMongo")
    public MongoTemplate towMongoTemplate() throws Exception {
        return new MongoTemplate(twoFactory(this.mongoProperties));
    }
 
    @Bean
    public SimpleMongoClientDatabaseFactory twoFactory(MongoProperties mongoProperties) throws Exception {
        return new SimpleMongoClientDatabaseFactory(mongoProperties.getUri());
    }
}
 
package com.ct.sync.config;
 
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.mongo.MongoProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.SimpleMongoClientDatabaseFactory;
import org.springframework.data.mongodb.repository.config.EnableMongoRepositories;

@Configuration
@EnableMongoRepositories(
        basePackages = "com.ct.sync.domain.*",
        mongoTemplateRef = "threeMongo")
public class ThreeMongoTemplate {
 
    @Autowired
    @Qualifier("threeMongoProperties")
    private MongoProperties mongoProperties;
 
    @Bean(name = "threeMongo")
    public MongoTemplate threeMongoTemplate() throws Exception {
        return new MongoTemplate(threeFactory(this.mongoProperties));
    }
 
    @Bean
    public SimpleMongoClientDatabaseFactory threeFactory(MongoProperties mongoProperties) throws Exception {
        return new SimpleMongoClientDatabaseFactory(mongoProperties.getUri());
    }
}
 
com.ct.sync.domain.*为mongodb实体类所在位置

3.创建mongo连接配置类

package com.ct.sync.config;
 
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.mongo.MongoProperties;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;
 
@Slf4j
@Configuration
public class MongoInit {

    @Bean(name = "oneMongoProperties")
    @Primary
    @ConfigurationProperties(prefix = "spring.data.mongodb.one")
    public MongoProperties oneMongoProperties() {
        log.info("-------------------- oneMongoProperties init ---------------------");
        return new MongoProperties();
    }
 
    @Bean(name = "twoMongoProperties")
    @ConfigurationProperties(prefix = "spring.data.mongodb.two")
    public MongoProperties twoMongoProperties() {
        log.info("-------------------- twoMongoProperties init ---------------------");
        return new MongoProperties();
    }
 
    @Bean(name = "threeMongoProperties")
    @ConfigurationProperties(prefix = "spring.data.mongodb.three")
    public MongoProperties threeMongoProperties() {
        log.info("-------------------- threeMongoProperties init ---------------------");
        return new MongoProperties();
    }
}

4.application.yml配置数据库连接

spring:
  data:
    mongodb:
      one:
        uri: mongodb://admin:123456@192.168.1.11:2332/DB_01
      two:
        uri: mongodb://admin:123456@192.168.1.11:2332/DB_02
      three:
        uri: mongodb://admin:123456@192.168.1.12:2332/DB_03

5.在指定的com.ct.sync.domain.*下创建实体类

package com.ct.sync.domain;

import lombok.Data;
import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;

@Data
@Document(collection = "record")
public class Record {

    @Id
    private String id; // ObjectId对应的字符串
    private String httpAddr;
    private String deviceId;
}

6.实现mongodb多数据源连接(我在项目中用定时任务同步数据)

package com.ct.sync.task;

import com.alibaba.fastjson.JSONObject;
import com.ct.sync.service.MiddleService;
import com.ct.sync.service.MongoSyncService;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.LocalTime;
import java.time.format.DateTimeFormatter;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

@Component
@EnableScheduling
@AllArgsConstructor
@Slf4j
public class MongoSyncTask {

    private final MongoSyncService mongoSyncService;
    
    @Autowired
    @Qualifier(value = "oneMongo")
    private MongoTemplate oneMongo;

    @Autowired
    @Qualifier(value = "twoMongo")
    private MongoTemplate twoMongo;

    @Autowired
    @Qualifier(value = "threeMongo")
    private MongoTemplate threeMongo;

    @Scheduled(cron = "0 30 11 * * *")
    public void runAllSync(){
        log.info("sync同步开启");
        ExecutorService executorService = Executors.newFixedThreadPool(3);
        executorService.submit(() -> mongoSyncService.runSync(oneMongo, "oneMongo"));
        executorService.submit(() -> mongoSyncService.runSync(twoMongo, "twoMongo"));
        executorService.submit(() -> mongoSyncService.runSync(threeMongo, "threeMongo"));
        executorService.shutdown();
        log.info("sync同步关闭");
    }
}
package com.ct.sync.service;

import org.springframework.data.mongodb.core.MongoTemplate;

public interface MongoSyncService {

    void runSync(MongoTemplate mongoTemplate, String logTag);
}
package com.ct.sync.service.impl;

import com.ct.sync.domain.Record;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

@Slf4j
@Service
public class MongoSyncServiceImpl implements MongoSyncService{

    @Transactional(rollbackFor = Exception.class)
    @Override
    public void runSync(MongoTemplate mongoTemplate, String logTag) {
        Query query = new Query();
       query.addCriteria(Criteria.where("startTime").gt(startOfDay.toInstant().toEpochMilli()));
        query.addCriteria(Criteria.where("endTime").lt(endOfDay.toInstant().toEpochMilli()));
        query.addCriteria(Criteria.where("deviceId").in(deviceCode));

        List<Record> records = mongoTemplate.find(query, Record.class);
         .....以下业务逻辑省略
    }
}

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐