springboot如何连接多个mongodb数据源
·
在实际工作中,我遇到了这种连接多个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);
.....以下业务逻辑省略
}
}
DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。
更多推荐


所有评论(0)