.ExpressionInterceptUrlRegistry registry = httpSecurity.authorizeRequests();
+ permitAllUrl.getUrls().forEach(url -> registry.antMatchers(url).permitAll());
+
+ httpSecurity
+ // CSRF禁用,因为不使用session
+ .csrf().disable()
+ .cors().configurationSource(corsConfigurationSource()).and()
+ // 禁用HTTP响应标头
+ .headers().cacheControl().disable().and()
+ // 认证失败处理类
+ .exceptionHandling().authenticationEntryPoint(unauthorizedHandler).and()
+ // 基于token,所以不需要session
+ .sessionManagement().sessionCreationPolicy(SessionCreationPolicy.STATELESS).and()
+ .sessionManagement().sessionFixation().none().and()
+ // 过滤请求
+ .authorizeRequests()
+ // 对于登录login 注册register 验证码captchaImage 允许匿名访问
+ .antMatchers("/api/thermal/**","/iotDA/device/**", "/websocket/allMemory", "/iot/UAV/**", "/iot/api/**", "/iot/device/userDevice", "/external/sms/sendAliSmsCode", "/login", "/register", "/captchaImage", "/iot/tool/register", "/iot/tool/ntp",
+ "/iot/tool/mqtt/auth", "/iot/workRecord/**", "/iot/tool/mqtt/authv5", "/iot/tool/mqtt/webhook", "/iot/tool/mqtt/webhookv5", "/auth/**/**",
+ "/wechat/mobileLogin", "/wechat/miniLogin", "/wechat/wxBind/callback").permitAll()
+ .antMatchers("/zlmhook/**").permitAll()
+ .antMatchers("/srs/**").permitAll() //srs服务
+ //.antMatchers("/sip/player/getBigScreenUrl/**").permitAll()
+ .antMatchers("/ruleengine/rulemanager/**").permitAll()
+ .antMatchers("/goview/sys/login", "/goview/project/getData").permitAll()
+ .antMatchers("/notify/smsLoginCaptcha", "/notify/emailRegisterCaptcha", "/auth/sms/login", "/notify/weComVerifyUrl"
+ , "/wechat/publicAccount/callback", "/notify/smsRegisterCaptcha").permitAll()
+ .antMatchers("/app/language/list").permitAll()
+ .antMatchers("/iot/news/bannerList", "/iot/news/topList", "/iot/news/getDetail").permitAll()
+ // 静态资源,可匿名访问
+ .antMatchers(HttpMethod.GET, "/", "/*.html", "/**/*.html", "/**/*.css", "/**/*.js", "/profile/**").permitAll()
+ .antMatchers("/swagger-ui.html", "/swagger-resources/**", "/webjars/**", "/*/api-docs", "/druid/**").permitAll()
+ .antMatchers("/druid/**").permitAll()
+ .antMatchers("/oauth2/**").permitAll()
+ .antMatchers("/swagger**", "/dev-api/swagger**", "/dev-api/swagger-ui/**", "/dev-api/v3/api-docs/**").permitAll()
+ .antMatchers("/**/device/**/**").permitAll()
+// // oauth
+// .antMatchers("/oauth/css/**","/oauth/fonts/**","/oauth/js/**").permitAll()
+ // dueros
+ .antMatchers("/dueros").permitAll()
+
+ .antMatchers(HttpMethod.OPTIONS, "/**").permitAll()
+
+ // 除上面外的所有请求全部需要鉴权认证
+ .anyRequest().authenticated()
+
+// // oauth
+// .and()
+// .formLogin()
+// .loginPage("/oauth/login")
+// .permitAll()
+// .and()
+// .logout().logoutUrl("/oauth/logout")
+// .permitAll()
+
+ .and()
+ .headers().frameOptions().disable();
+ // 添加Logout filter
+ httpSecurity.logout().logoutUrl("/logout").logoutSuccessHandler(logoutSuccessHandler);
+ // 添加JWT filter
+ httpSecurity.addFilterBefore(authenticationTokenFilter, UsernamePasswordAuthenticationFilter.class);
+ // 添加CORS filter
+ httpSecurity.addFilterBefore(corsFilter, JwtAuthenticationTokenFilter.class);
+ httpSecurity.addFilterBefore(corsFilter, LogoutFilter.class);
+ }
+
+ /**
+ * 强散列哈希加密实现
+ */
+ @Bean
+ public BCryptPasswordEncoder bCryptPasswordEncoder() {
+ return new BCryptPasswordEncoder();
+ }
+
+ /**
+ * 身份认证接口
+ */
+ @Override
+ protected void configure(AuthenticationManagerBuilder auth) throws Exception {
+ auth.userDetailsService(userDetailsService).passwordEncoder(bCryptPasswordEncoder());
+ }
+
+ @Bean
+ public CorsConfigurationSource corsConfigurationSource() {
+// CorsConfiguration config = new CorsConfiguration();
+// config.setAllowCredentials(true);
+// config.addAllowedOriginPattern("*"); // 通配支持
+// config.addAllowedHeader("*");
+// config.addAllowedMethod("*");
+// config.setMaxAge(1800L);
+//
+// UrlBasedCorsConfigurationSource source = new UrlBasedCorsConfigurationSource();
+// source.registerCorsConfiguration("/**", config);
+// return source;
+
+
+ CorsConfiguration config = new CorsConfiguration();
+ // 关键:设置允许的前端源(替换为你的实际前端域名,如https://ri.satabot.com)
+ // 填* 会有问题
+ config.setAllowedOrigins(Arrays.asList("https://ri.satabot.com", "https://ai.satabot.com",
+ "https://serviceai.satabot.com", "https://serviceri.satabot.com", "https://pvs.satabot.com",
+ "https://servicepvs.satabot.com", "https://servicepathplan.satabot.com", "https://solar.satabot.com",
+ "https://servicestream.satabot.com", "http://192.168.2.58", "http://192.168.2.58"));
+// config.addAllowedOriginPattern("*");
+ // 允许的请求方法(必须包含OPTIONS,支持业务请求的GET/POST等)
+ config.setAllowedMethods(Arrays.asList("GET", "POST", "PUT", "DELETE", "OPTIONS"));
+ // 允许的请求头(*表示所有,若有自定义头可单独指定,如Token)
+ config.setAllowedHeaders(Arrays.asList("*"));
+ // 允许携带Cookie/Token(开启后,allowedOrigins不能用*,需指定具体域名)
+ config.setAllowCredentials(true);
+ // 预检请求缓存时间(单位:秒),减少OPTIONS请求次数
+ config.setMaxAge(3600L);
+
+ // 将跨域配置绑定到所有接口路径
+ UrlBasedCorsConfigurationSource source = new UrlBasedCorsConfigurationSource();
+ source.registerCorsConfiguration("/**", config);
+ return source;
+ }
+
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/ServerConfig.java b/maibu-framework/src/main/java/com/maibu/config/ServerConfig.java
new file mode 100644
index 0000000..540b437
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/ServerConfig.java
@@ -0,0 +1,32 @@
+package com.maibu.config;
+
+import javax.servlet.http.HttpServletRequest;
+import org.springframework.stereotype.Component;
+import com.fastbee.common.utils.ServletUtils;
+
+/**
+ * 服务相关配置
+ *
+ * @author ruoyi
+ */
+@Component
+public class ServerConfig
+{
+ /**
+ * 获取完整的请求路径,包括:域名,端口,上下文访问路径
+ *
+ * @return 服务地址
+ */
+ public String getUrl()
+ {
+ HttpServletRequest request = ServletUtils.getRequest();
+ return getDomain(request);
+ }
+
+ public static String getDomain(HttpServletRequest request)
+ {
+ StringBuffer url = request.getRequestURL();
+ String contextPath = request.getServletContext().getContextPath();
+ return url.delete(url.length() - request.getRequestURI().length(), url.length()).append(contextPath).toString();
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/SqlFilterArgumentResolver.java b/maibu-framework/src/main/java/com/maibu/config/SqlFilterArgumentResolver.java
new file mode 100644
index 0000000..939e0cc
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/SqlFilterArgumentResolver.java
@@ -0,0 +1,92 @@
+package com.maibu.config;
+
+import cn.hutool.core.util.StrUtil;
+import com.baomidou.mybatisplus.core.metadata.OrderItem;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.core.MethodParameter;
+import org.springframework.web.bind.support.WebDataBinderFactory;
+import org.springframework.web.context.request.NativeWebRequest;
+import org.springframework.web.method.support.HandlerMethodArgumentResolver;
+import org.springframework.web.method.support.ModelAndViewContainer;
+
+import javax.servlet.http.HttpServletRequest;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Optional;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+/**
+ * 解决Mybatis Plus Order By SQL注入问题
+ * @author admin
+ */
+@Slf4j
+public class SqlFilterArgumentResolver implements HandlerMethodArgumentResolver {
+
+ private final static String[] KEYWORDS = { "master", "truncate", "insert", "select", "delete", "update", "declare",
+ "alter", "drop", "sleep" };
+
+ /**
+ * 判断Controller是否包含page 参数
+ * @param parameter 参数
+ * @return 是否过滤
+ */
+ @Override
+ public boolean supportsParameter(MethodParameter parameter) {
+ return parameter.getParameterType().equals(Page.class);
+ }
+
+ /**
+ * @param parameter 入参集合
+ * @param mavContainer model 和 view
+ * @param webRequest web相关
+ * @param binderFactory 入参解析
+ * @return 检查后新的page对象
+ *
+ * page 只支持查询 GET .如需解析POST获取请求报文体处理
+ */
+ @Override
+ public Object resolveArgument(MethodParameter parameter, ModelAndViewContainer mavContainer,
+ NativeWebRequest webRequest, WebDataBinderFactory binderFactory) {
+
+ HttpServletRequest request = webRequest.getNativeRequest(HttpServletRequest.class);
+
+ String[] ascs = request.getParameterValues("ascs");
+ String[] descs = request.getParameterValues("descs");
+ String current = request.getParameter("current");
+ String size = request.getParameter("size");
+
+ Page> page = new Page<>();
+ if (StrUtil.isNotBlank(current)) {
+ page.setCurrent(Long.parseLong(current));
+ }
+
+ if (StrUtil.isNotBlank(size)) {
+ page.setSize(Long.parseLong(size));
+ }
+ List orderItemList = new ArrayList<>();
+ Optional.ofNullable(ascs).ifPresent(s -> orderItemList.addAll(
+ Arrays.stream(s).filter(sqlInjectPredicate()).map(OrderItem::asc).collect(Collectors.toList())));
+ Optional.ofNullable(descs).ifPresent(s -> orderItemList.addAll(
+ Arrays.stream(s).filter(sqlInjectPredicate()).map(OrderItem::desc).collect(Collectors.toList())));
+ page.addOrder(orderItemList);
+ return page;
+ }
+
+ /**
+ * 判断用户输入里面有没有关键字
+ * @return Predicate
+ */
+ private Predicate sqlInjectPredicate() {
+ return sql -> {
+ for (String keyword : KEYWORDS) {
+ if (StrUtil.containsIgnoreCase(sql, keyword)) {
+ return false;
+ }
+ }
+ return true;
+ };
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/ThreadPoolConfig.java b/maibu-framework/src/main/java/com/maibu/config/ThreadPoolConfig.java
new file mode 100644
index 0000000..bcaa998
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/ThreadPoolConfig.java
@@ -0,0 +1,63 @@
+package com.maibu.config;
+
+import com.fastbee.common.utils.Threads;
+import org.apache.commons.lang3.concurrent.BasicThreadFactory;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.ThreadPoolExecutor;
+
+/**
+ * 线程池配置
+ *
+ * @author ruoyi
+ **/
+@Configuration
+public class ThreadPoolConfig
+{
+ // 核心线程池大小
+ private int corePoolSize = 50;
+
+ // 最大可创建的线程数
+ private int maxPoolSize = 200;
+
+ // 队列最大长度
+ private int queueCapacity = 1000;
+
+ // 线程池维护线程所允许的空闲时间
+ private int keepAliveSeconds = 300;
+
+ @Bean(name = "threadPoolTaskExecutor")
+ public ThreadPoolTaskExecutor threadPoolTaskExecutor()
+ {
+ ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
+ executor.setMaxPoolSize(maxPoolSize);
+ executor.setCorePoolSize(corePoolSize);
+ executor.setQueueCapacity(queueCapacity);
+ executor.setKeepAliveSeconds(keepAliveSeconds);
+ // 线程池对拒绝任务(无线程可用)的处理策略
+ executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
+ return executor;
+ }
+
+ /**
+ * 执行周期性或定时任务
+ */
+ @Bean(name = "scheduledExecutorService")
+ protected ScheduledExecutorService scheduledExecutorService()
+ {
+ return new ScheduledThreadPoolExecutor(corePoolSize,
+ new BasicThreadFactory.Builder().namingPattern("schedule-pool-%d").daemon(true).build(),
+ new ThreadPoolExecutor.CallerRunsPolicy())
+ {
+ @Override
+ protected void afterExecute(Runnable r, Throwable t)
+ {
+ super.afterExecute(r, t);
+ Threads.printException(r, t);
+ }
+ };
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/properties/DruidProperties.java b/maibu-framework/src/main/java/com/maibu/config/properties/DruidProperties.java
new file mode 100644
index 0000000..ee9e961
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/properties/DruidProperties.java
@@ -0,0 +1,75 @@
+package com.maibu.config.properties;
+
+import org.springframework.context.annotation.Configuration;
+
+/**
+ * druid 配置属性
+ *
+ * @author ruoyi
+ */
+@Configuration
+public class DruidProperties
+{
+// @Value("${spring.datasource.druid.initialSize}")
+// private int initialSize;
+//
+// @Value("${spring.datasource.druid.minIdle}")
+// private int minIdle;
+//
+// @Value("${spring.datasource.druid.maxActive}")
+// private int maxActive;
+//
+// @Value("${spring.datasource.druid.maxWait}")
+// private int maxWait;
+//
+// @Value("${spring.datasource.druid.timeBetweenEvictionRunsMillis}")
+// private int timeBetweenEvictionRunsMillis;
+//
+// @Value("${spring.datasource.druid.minEvictableIdleTimeMillis}")
+// private int minEvictableIdleTimeMillis;
+//
+// @Value("${spring.datasource.druid.maxEvictableIdleTimeMillis}")
+// private int maxEvictableIdleTimeMillis;
+//
+// @Value("${spring.datasource.druid.validationQuery}")
+// private String validationQuery;
+//
+// @Value("${spring.datasource.druid.testWhileIdle}")
+// private boolean testWhileIdle;
+//
+// @Value("${spring.datasource.druid.testOnBorrow}")
+// private boolean testOnBorrow;
+//
+// @Value("${spring.datasource.druid.testOnReturn}")
+// private boolean testOnReturn;
+
+// public DruidDataSource dataSource(DruidDataSource datasource)
+// {
+// /** 配置初始化大小、最小、最大 */
+// datasource.setInitialSize(initialSize);
+// datasource.setMaxActive(maxActive);
+// datasource.setMinIdle(minIdle);
+//
+// /** 配置获取连接等待超时的时间 */
+// datasource.setMaxWait(maxWait);
+//
+// /** 配置间隔多久才进行一次检测,检测需要关闭的空闲连接,单位是毫秒 */
+// datasource.setTimeBetweenEvictionRunsMillis(timeBetweenEvictionRunsMillis);
+//
+// /** 配置一个连接在池中最小、最大生存的时间,单位是毫秒 */
+// datasource.setMinEvictableIdleTimeMillis(minEvictableIdleTimeMillis);
+// datasource.setMaxEvictableIdleTimeMillis(maxEvictableIdleTimeMillis);
+//
+// /**
+// * 用来检测连接是否有效的sql,要求是一个查询语句,常用select 'x'。如果validationQuery为null,testOnBorrow、testOnReturn、testWhileIdle都不会起作用。
+// */
+// datasource.setValidationQuery(validationQuery);
+// /** 建议配置为true,不影响性能,并且保证安全性。申请连接的时候检测,如果空闲时间大于timeBetweenEvictionRunsMillis,执行validationQuery检测连接是否有效。 */
+// datasource.setTestWhileIdle(testWhileIdle);
+// /** 申请连接时执行validationQuery检测连接是否有效,做了这个配置会降低性能。 */
+// datasource.setTestOnBorrow(testOnBorrow);
+// /** 归还连接时执行validationQuery检测连接是否有效,做了这个配置会降低性能。 */
+// datasource.setTestOnReturn(testOnReturn);
+// return datasource;
+// }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/properties/PermitAllUrlProperties.java b/maibu-framework/src/main/java/com/maibu/config/properties/PermitAllUrlProperties.java
new file mode 100644
index 0000000..6edbca3
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/properties/PermitAllUrlProperties.java
@@ -0,0 +1,72 @@
+package com.maibu.config.properties;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.regex.Pattern;
+import org.apache.commons.lang3.RegExUtils;
+import org.springframework.beans.BeansException;
+import org.springframework.beans.factory.InitializingBean;
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.ApplicationContextAware;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.core.annotation.AnnotationUtils;
+import org.springframework.web.method.HandlerMethod;
+import org.springframework.web.servlet.mvc.method.RequestMappingInfo;
+import org.springframework.web.servlet.mvc.method.annotation.RequestMappingHandlerMapping;
+import com.fastbee.common.annotation.Anonymous;
+
+/**
+ * 设置Anonymous注解允许匿名访问的url
+ *
+ * @author ruoyi
+ */
+@Configuration
+public class PermitAllUrlProperties implements InitializingBean, ApplicationContextAware
+{
+ private static final Pattern PATTERN = Pattern.compile("\\{(.*?)\\}");
+
+ private ApplicationContext applicationContext;
+
+ private List urls = new ArrayList<>();
+
+ public String ASTERISK = "*";
+
+ @Override
+ public void afterPropertiesSet()
+ {
+ RequestMappingHandlerMapping mapping = applicationContext.getBean(RequestMappingHandlerMapping.class);
+ Map map = mapping.getHandlerMethods();
+
+ map.keySet().forEach(info -> {
+ HandlerMethod handlerMethod = map.get(info);
+
+ // 获取方法上边的注解 替代path variable 为 *
+ Anonymous method = AnnotationUtils.findAnnotation(handlerMethod.getMethod(), Anonymous.class);
+ Optional.ofNullable(method).ifPresent(anonymous -> info.getPatternsCondition().getPatterns()
+ .forEach(url -> urls.add(RegExUtils.replaceAll(url, PATTERN, ASTERISK))));
+
+ // 获取类上边的注解, 替代path variable 为 *
+ Anonymous controller = AnnotationUtils.findAnnotation(handlerMethod.getBeanType(), Anonymous.class);
+ Optional.ofNullable(controller).ifPresent(anonymous -> info.getPatternsCondition().getPatterns()
+ .forEach(url -> urls.add(RegExUtils.replaceAll(url, PATTERN, ASTERISK))));
+ });
+ }
+
+ @Override
+ public void setApplicationContext(ApplicationContext context) throws BeansException
+ {
+ this.applicationContext = context;
+ }
+
+ public List getUrls()
+ {
+ return urls;
+ }
+
+ public void setUrls(List urls)
+ {
+ this.urls = urls;
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/properties/RedissonProperties.java b/maibu-framework/src/main/java/com/maibu/config/properties/RedissonProperties.java
new file mode 100644
index 0000000..3ae6fce
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/properties/RedissonProperties.java
@@ -0,0 +1,134 @@
+package com.maibu.config.properties;
+
+import lombok.Data;
+import lombok.NoArgsConstructor;
+import org.redisson.config.ReadMode;
+import org.redisson.config.SubscriptionMode;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+
+/**
+ * Redisson 配置属性
+ *
+ */
+@Data
+@ConfigurationProperties(prefix = "redisson")
+public class RedissonProperties {
+
+ /**
+ * redis缓存key前缀
+ */
+ private String keyPrefix;
+
+ /**
+ * 线程池数量,默认值 = 当前处理核数量 * 2
+ */
+ private int threads;
+
+ /**
+ * Netty线程池数量,默认值 = 当前处理核数量 * 2
+ */
+ private int nettyThreads;
+
+ /**
+ * 单机服务配置
+ */
+ private SingleServerConfig singleServerConfig;
+
+ /**
+ * 集群服务配置
+ */
+ private ClusterServersConfig clusterServersConfig;
+
+ @Data
+ @NoArgsConstructor
+ public static class SingleServerConfig {
+
+ /**
+ * 客户端名称
+ */
+ private String clientName;
+
+ /**
+ * 最小空闲连接数
+ */
+ private int connectionMinimumIdleSize;
+
+ /**
+ * 连接池大小
+ */
+ private int connectionPoolSize;
+
+ /**
+ * 连接空闲超时,单位:毫秒
+ */
+ private int idleConnectionTimeout;
+
+ /**
+ * 命令等待超时,单位:毫秒
+ */
+ private int timeout;
+
+ /**
+ * 发布和订阅连接池大小
+ */
+ private int subscriptionConnectionPoolSize;
+
+ }
+
+ @Data
+ @NoArgsConstructor
+ public static class ClusterServersConfig {
+
+ /**
+ * 客户端名称
+ */
+ private String clientName;
+
+ /**
+ * master最小空闲连接数
+ */
+ private int masterConnectionMinimumIdleSize;
+
+ /**
+ * master连接池大小
+ */
+ private int masterConnectionPoolSize;
+
+ /**
+ * slave最小空闲连接数
+ */
+ private int slaveConnectionMinimumIdleSize;
+
+ /**
+ * slave连接池大小
+ */
+ private int slaveConnectionPoolSize;
+
+ /**
+ * 连接空闲超时,单位:毫秒
+ */
+ private int idleConnectionTimeout;
+
+ /**
+ * 命令等待超时,单位:毫秒
+ */
+ private int timeout;
+
+ /**
+ * 发布和订阅连接池大小
+ */
+ private int subscriptionConnectionPoolSize;
+
+ /**
+ * 读取模式
+ */
+ private ReadMode readMode;
+
+ /**
+ * 订阅模式
+ */
+ private SubscriptionMode subscriptionMode;
+
+ }
+
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingAlgorithmTool.java b/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingAlgorithmTool.java
new file mode 100644
index 0000000..ac819a3
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingAlgorithmTool.java
@@ -0,0 +1,262 @@
+package com.maibu.config.sharding;
+
+import cn.hutool.extra.spring.SpringUtil;
+import com.alibaba.druid.util.StringUtils;
+import com.baomidou.mybatisplus.core.toolkit.CollectionUtils;
+import com.fastbee.framework.config.sharding.enums.ShardingTableCacheEnum;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.shardingsphere.driver.jdbc.core.datasource.ShardingSphereDataSource;
+import org.apache.shardingsphere.infra.config.RuleConfiguration;
+import org.apache.shardingsphere.mode.manager.ContextManager;
+import org.apache.shardingsphere.sharding.algorithm.config.AlgorithmProvidedShardingRuleConfiguration;
+import org.apache.shardingsphere.sharding.api.config.rule.ShardingTableRuleConfiguration;
+import org.springframework.core.env.Environment;
+
+import java.sql.*;
+import java.time.YearMonth;
+import java.time.format.DateTimeFormatter;
+import java.util.*;
+import java.util.stream.Collectors;
+
+/**
+ * @Title ShardingAlgorithmTool
+ *
@Description 按月分片算法工具
+ *
+ */
+@Slf4j
+public class ShardingAlgorithmTool {
+
+ /** 表分片符号,例:siot_device_log_202201 中,分片符号为 "_" */
+ private static final String TABLE_SPLIT_SYMBOL = "_";
+
+ /** 数据库配置 */
+ private static final Environment ENV = SpringUtil.getApplicationContext().getEnvironment();
+ private static final String DATASOURCE_URL = ENV.getProperty("spring.shardingsphere.datasource.ds0.url");
+ private static final String DATASOURCE_USERNAME = ENV.getProperty("spring.shardingsphere.datasource.ds0.username");
+ private static final String DATASOURCE_PASSWORD = ENV.getProperty("spring.shardingsphere.datasource.ds0.password");
+
+
+ /**
+ * 检查分表获取的表名是否存在,不存在则自动建表
+ * @param logicTable 逻辑表
+ * @param resultTableNames 真实表名,例:iot_device_log_202201
+ * @return 存在于数据库中的真实表名集合
+ */
+ public static Set getShardingTablesAndCreate(ShardingTableCacheEnum logicTable, Collection resultTableNames) {
+ return resultTableNames.stream().map(o -> getShardingTableAndCreate(logicTable, o)).collect(Collectors.toSet());
+ }
+
+ /**
+ * 检查分表获取的表名是否存在,不存在则自动建表
+ * @param logicTable 逻辑表
+ * @param resultTableName 真实表名,例:iot_device_log_202201
+ * @return 确认存在于数据库中的真实表名
+ */
+ public static String getShardingTableAndCreate(ShardingTableCacheEnum logicTable, String resultTableName) {
+ // 缓存中有此表则返回,没有则判断创建
+ if (logicTable.resultTableNamesCache().contains(resultTableName)) {
+ return resultTableName;
+ } else {
+ // 未创建的表返回逻辑空表
+ boolean isSuccess = createShardingTable(logicTable, resultTableName);
+ return isSuccess ? resultTableName : logicTable.logicTableName();
+ }
+ }
+
+ /**
+ * 重载全部缓存
+ */
+ public static void tableNameCacheReloadAll() {
+ Arrays.stream(ShardingTableCacheEnum.values()).forEach(ShardingAlgorithmTool::tableNameCacheReload);
+ }
+
+ /**
+ * 重载指定分表缓存
+ * @param logicTable 逻辑表
+ */
+ public static void tableNameCacheReload(ShardingTableCacheEnum logicTable) {
+ // 读取数据库中所有表名
+ List tableNameList = getAllTableNameBySchema(logicTable);
+ // 更新缓存、配置(原子操作)
+ logicTable.atomicUpdateCacheAndActualDataNodes(tableNameList);
+ // 删除旧的缓存(如果存在)
+ logicTable.resultTableNamesCache().clear();
+ // 写入新的缓存
+ logicTable.resultTableNamesCache().addAll(tableNameList);
+ // 动态更新配置 actualDataNodes
+ actualDataNodesRefresh(logicTable.logicTableName(), tableNameList);
+ }
+
+ /**
+ * 获取所有表名
+ * @return 表名集合
+ * @param logicTable 逻辑表
+ */
+ public static List getAllTableNameBySchema(ShardingTableCacheEnum logicTable) {
+ List tableNames = new ArrayList<>();
+ if (StringUtils.isEmpty(DATASOURCE_URL) || StringUtils.isEmpty(DATASOURCE_USERNAME) || StringUtils.isEmpty(DATASOURCE_PASSWORD)) {
+ log.error(">>>>>>>>>> 【ERROR】数据库连接配置有误,请稍后重试,URL:{}, username:{}, password:{}", DATASOURCE_URL, DATASOURCE_USERNAME, DATASOURCE_PASSWORD);
+ throw new IllegalArgumentException("数据库连接配置有误,请稍后重试");
+ }
+ try (Connection conn = DriverManager.getConnection(DATASOURCE_URL, DATASOURCE_USERNAME, DATASOURCE_PASSWORD);
+ Statement st = conn.createStatement()) {
+ String logicTableName = logicTable.logicTableName();
+ try (ResultSet rs = st.executeQuery("show TABLES like '" + logicTableName + TABLE_SPLIT_SYMBOL + "%'")) {
+ log.info("查询数据库所有表:{}","show TABLES like '" + logicTableName + TABLE_SPLIT_SYMBOL + "%'");
+ while (rs.next()) {
+ String tableName = rs.getString(1);
+ log.info("分表格式:{}",String.format("^(%s\\d{6})$", logicTableName + TABLE_SPLIT_SYMBOL));
+ // 匹配分表格式 例:^(t\_contract_\d{6})$
+ if (org.apache.commons.lang3.StringUtils.isNotBlank(tableName) && tableName.matches(String.format("^(%s\\d{6})$", logicTableName + TABLE_SPLIT_SYMBOL))) {
+ tableNames.add(rs.getString(1));
+ }
+ }
+ }
+ } catch (SQLException e) {
+ log.error(">>>>>>>>>> 【ERROR】数据库连接失败,请稍后重试,原因:{}", e.getMessage(), e);
+ throw new IllegalArgumentException("数据库连接失败,请稍后重试");
+ }
+ return tableNames;
+ }
+
+ /**
+ * 动态更新配置 actualDataNodes
+ *
+ * @param logicTableName 逻辑表名
+ * @param tableNamesCache 真实表名集合
+ */
+ public static void actualDataNodesRefresh(String logicTableName, List tableNamesCache) {
+ try {
+ if (CollectionUtils.isEmpty(tableNamesCache)) {
+ return;
+ }
+ // 获取数据分片节点
+ String dbName = "ds0";
+ log.info(">>>>>>>>>> 【INFO】更新分表配置,logicTableName:{},tableNamesCache:{}", logicTableName, tableNamesCache);
+
+ // generate actualDataNodes
+ String newActualDataNodes = tableNamesCache.stream().map(o -> String.format("%s.%s", dbName, o)).collect(Collectors.joining(","));
+ ShardingSphereDataSource shardingSphereDataSource = SpringUtil.getBean(ShardingSphereDataSource.class);
+ updateShardRuleActualDataNodes(shardingSphereDataSource, logicTableName, newActualDataNodes);
+ }catch (Exception e){
+ log.error("初始化 动态表单失败,原因:{}", e.getMessage(), e);
+ }
+ }
+
+
+ // --------------------------------------------------------------------------------------------------------------
+ // 私有方法
+ // --------------------------------------------------------------------------------------------------------------
+
+
+ /**
+ * 刷新ActualDataNodes
+ */
+ private static void updateShardRuleActualDataNodes(ShardingSphereDataSource dataSource, String logicTableName, String newActualDataNodes) {
+ // Context manager.
+ ContextManager contextManager = dataSource.getContextManager();
+
+ // Rule configuration.
+ String schemaName = "logic_db";
+ Collection newRuleConfigList = new LinkedList<>();
+ Collection oldRuleConfigList = dataSource.getContextManager()
+ .getMetaDataContexts()
+ .getMetaData(schemaName)
+ .getRuleMetaData()
+ .getConfigurations();
+
+ for (RuleConfiguration oldRuleConfig : oldRuleConfigList) {
+ if (oldRuleConfig instanceof AlgorithmProvidedShardingRuleConfiguration) {
+
+ // Algorithm provided sharding rule configuration
+ AlgorithmProvidedShardingRuleConfiguration oldAlgorithmConfig = (AlgorithmProvidedShardingRuleConfiguration) oldRuleConfig;
+ AlgorithmProvidedShardingRuleConfiguration newAlgorithmConfig = new AlgorithmProvidedShardingRuleConfiguration();
+
+ // Sharding table rule configuration Collection
+ Collection newTableRuleConfigList = new LinkedList<>();
+ Collection oldTableRuleConfigList = oldAlgorithmConfig.getTables();
+
+ oldTableRuleConfigList.forEach(oldTableRuleConfig -> {
+ if (logicTableName.equals(oldTableRuleConfig.getLogicTable())) {
+ ShardingTableRuleConfiguration newTableRuleConfig = new ShardingTableRuleConfiguration(oldTableRuleConfig.getLogicTable(), newActualDataNodes);
+ newTableRuleConfig.setTableShardingStrategy(oldTableRuleConfig.getTableShardingStrategy());
+ newTableRuleConfig.setDatabaseShardingStrategy(oldTableRuleConfig.getDatabaseShardingStrategy());
+ newTableRuleConfig.setKeyGenerateStrategy(oldTableRuleConfig.getKeyGenerateStrategy());
+
+ newTableRuleConfigList.add(newTableRuleConfig);
+ } else {
+ newTableRuleConfigList.add(oldTableRuleConfig);
+ }
+ });
+
+ newAlgorithmConfig.setTables(newTableRuleConfigList);
+ newAlgorithmConfig.setAutoTables(oldAlgorithmConfig.getAutoTables());
+ newAlgorithmConfig.setBindingTableGroups(oldAlgorithmConfig.getBindingTableGroups());
+ newAlgorithmConfig.setBroadcastTables(oldAlgorithmConfig.getBroadcastTables());
+ newAlgorithmConfig.setDefaultDatabaseShardingStrategy(oldAlgorithmConfig.getDefaultDatabaseShardingStrategy());
+ newAlgorithmConfig.setDefaultTableShardingStrategy(oldAlgorithmConfig.getDefaultTableShardingStrategy());
+ newAlgorithmConfig.setDefaultKeyGenerateStrategy(oldAlgorithmConfig.getDefaultKeyGenerateStrategy());
+ newAlgorithmConfig.setDefaultShardingColumn(oldAlgorithmConfig.getDefaultShardingColumn());
+ newAlgorithmConfig.setShardingAlgorithms(oldAlgorithmConfig.getShardingAlgorithms());
+ newAlgorithmConfig.setKeyGenerators(oldAlgorithmConfig.getKeyGenerators());
+ newRuleConfigList.add(newAlgorithmConfig);
+ }
+ }
+
+ // update context
+ contextManager.alterRuleConfiguration(schemaName, newRuleConfigList);
+ }
+
+ /**
+ * 创建分表
+ * @param logicTable 逻辑表
+ * @param resultTableName 真实表名,例:sys_user_behavior_202201
+ * @return 创建结果(true创建成功,false未创建)
+ */
+ private static boolean createShardingTable(ShardingTableCacheEnum logicTable, String resultTableName) {
+ // 根据日期判断,当前月份之后分表不提前创建
+ String month = resultTableName.replace(logicTable.logicTableName() + TABLE_SPLIT_SYMBOL,"");
+ YearMonth shardingMonth = YearMonth.parse(month, DateTimeFormatter.ofPattern("yyyyMM"));
+ if (shardingMonth.isAfter(YearMonth.now())) {
+ return false;
+ }
+
+ synchronized (logicTable.logicTableName().intern()) {
+ // 缓存中有此表 返回
+ if (logicTable.resultTableNamesCache().contains(resultTableName)) {
+ return false;
+ }
+ // 缓存中无此表,则建表并添加缓存
+ executeSql(Collections.singletonList("CREATE TABLE IF NOT EXISTS `" + resultTableName + "` LIKE `" + logicTable.logicTableName() + "`;"));
+ // 缓存重载
+ tableNameCacheReload(logicTable);
+ }
+ return true;
+ }
+
+ /**
+ * 执行SQL
+ * @param sqlList SQL集合
+ */
+ private static void executeSql(List sqlList) {
+ if (StringUtils.isEmpty(DATASOURCE_URL) || StringUtils.isEmpty(DATASOURCE_USERNAME) || StringUtils.isEmpty(DATASOURCE_PASSWORD)) {
+ log.error(">>>>>>>>>> 【ERROR】数据库连接配置有误,请稍后重试,URL:{}, username:{}, password:{}", DATASOURCE_URL, DATASOURCE_USERNAME, DATASOURCE_PASSWORD);
+ throw new IllegalArgumentException("数据库连接配置有误,请稍后重试");
+ }
+ try (Connection conn = DriverManager.getConnection(DATASOURCE_URL, DATASOURCE_USERNAME, DATASOURCE_PASSWORD)) {
+ try (Statement st = conn.createStatement()) {
+ conn.setAutoCommit(false);
+ for (String sql : sqlList) {
+ st.execute(sql);
+ }
+ } catch (Exception e) {
+ conn.rollback();
+ log.error(">>>>>>>>>> 【ERROR】数据表创建执行失败,请稍后重试,原因:{}", e.getMessage(), e);
+ throw new IllegalArgumentException("数据表创建执行失败,请稍后重试");
+ }
+ } catch (SQLException e) {
+ log.error(">>>>>>>>>> 【ERROR】数据库连接失败,请稍后重试,原因:{}", e.getMessage(), e);
+ throw new IllegalArgumentException("数据库连接失败,请稍后重试");
+ }
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingTablesLoadRunner.java b/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingTablesLoadRunner.java
new file mode 100644
index 0000000..6c65231
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/sharding/ShardingTablesLoadRunner.java
@@ -0,0 +1,21 @@
+package com.maibu.config.sharding;
+
+import org.springframework.boot.CommandLineRunner;
+import org.springframework.core.annotation.Order;
+import org.springframework.stereotype.Component;
+
+/**
+ * @Title ShardingTablesLoadRunner
+ *
@Description 项目启动后,读取已有分表,进行缓存
+ *
+ */
+@Order(value = 1) // 数字越小,越先执行
+@Component
+public class ShardingTablesLoadRunner implements CommandLineRunner {
+
+ @Override
+ public void run(String... args) {
+ // 读取已有分表,进行缓存
+ //ShardingAlgorithmTool.tableNameCacheReloadAll();
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/sharding/TimeShardingAlgorithm.java b/maibu-framework/src/main/java/com/maibu/config/sharding/TimeShardingAlgorithm.java
new file mode 100644
index 0000000..89eea86
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/sharding/TimeShardingAlgorithm.java
@@ -0,0 +1,191 @@
+package com.maibu.config.sharding;
+
+import com.fastbee.framework.config.sharding.enums.ShardingTableCacheEnum;
+import com.google.common.collect.Range;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.shardingsphere.sharding.api.sharding.standard.PreciseShardingValue;
+import org.apache.shardingsphere.sharding.api.sharding.standard.RangeShardingValue;
+import org.apache.shardingsphere.sharding.api.sharding.standard.StandardShardingAlgorithm;
+import org.springframework.util.CollectionUtils;
+
+import java.text.SimpleDateFormat;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.time.ZonedDateTime;
+import java.time.format.DateTimeFormatter;
+import java.util.*;
+import java.util.function.Function;
+
+/**
+ *
@Title TimeShardingAlgorithm
+ *
@Description 分片算法,按月分片
+ *
+ */
+@Slf4j
+public class TimeShardingAlgorithm implements StandardShardingAlgorithm {
+ /**
+ * Date类型的分片时间格式
+ */
+ private static final SimpleDateFormat TABLE_SHARD_Date_FORMATTER = new SimpleDateFormat("yyyyMM");
+
+ /**
+ * 分片时间格式
+ */
+ private static final DateTimeFormatter TABLE_SHARD_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyyMM");
+
+ /**
+ * 完整时间格式
+ */
+ private static final DateTimeFormatter DATE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyyMMdd HH:mm:ss");
+ /**
+ * 完整时间格式
+ */
+ private static final SimpleDateFormat DATE_TIME_FORMATTER_SPILE = new SimpleDateFormat("yyyy-MM-dd");
+
+ /**
+ * 表分片符号,例:t_user_202201 中,分片符号为 "_"
+ */
+ private final String TABLE_SPLIT_SYMBOL = "_";
+
+
+ /**
+ * 精准分片
+ * @param tableNames 对应分片库中所有分片表的集合
+ * @param preciseShardingValue 分片键值,其中 logicTableName 为逻辑表,columnName 分片键,value 为从 SQL 中解析出来的分片键的值
+ * @return 表名
+ */
+ @Override
+ public String doSharding(Collection tableNames, PreciseShardingValue preciseShardingValue) {
+ String logicTableName = preciseShardingValue.getLogicTableName();
+ ShardingTableCacheEnum logicTable = ShardingTableCacheEnum.of(logicTableName);
+ createAllTable(logicTable, tableNames);
+
+ /// 打印分片信息
+ log.info(">>>>>>>>>> 【INFO】精确分片,节点配置表名:{},数据库缓存表名:{}", tableNames, logicTable.resultTableNamesCache());
+
+ Date date = preciseShardingValue.getValue();
+ Instant instant = date.toInstant();
+ LocalDateTime localDateTime = instant.atZone(ZoneId.systemDefault()).toLocalDateTime();
+ String resultTableName = logicTableName + "_" + TABLE_SHARD_TIME_FORMATTER.format(localDateTime);
+ // 检查分表获取的表名是否存在,不存在则自动建表
+ if (!tableNames.contains(resultTableName)){
+ tableNames.add(resultTableName);
+ }
+ return ShardingAlgorithmTool.getShardingTableAndCreate(logicTable, resultTableName);
+ }
+
+ /**
+ * 范围分片
+ * @param tableNames 对应分片库中所有分片表的集合
+ * @param rangeShardingValue 分片范围
+ * @return 表名集合
+ */
+ @Override
+ public Collection doSharding(Collection tableNames, RangeShardingValue rangeShardingValue) {
+ log.info("开始分表查询开始:{}",System.currentTimeMillis());
+ String logicTableName = rangeShardingValue.getLogicTableName();
+ ShardingTableCacheEnum logicTable = ShardingTableCacheEnum.of(logicTableName);
+ createAllTable(logicTable, tableNames);
+
+ /// 打印分片信息
+ log.info(">>>>>>>>>> 【INFO】范围分片,节点配置表名:{},数据库缓存表名:{}", tableNames, logicTable.resultTableNamesCache());
+
+ // between and 的起始值
+ Range valueRange = rangeShardingValue.getValueRange();
+ boolean hasLowerBound = valueRange.hasLowerBound();
+ boolean hasUpperBound = valueRange.hasUpperBound();
+
+ // 获取最大值和最小值
+ Set tableNameCache = logicTable.resultTableNamesCache();
+ String min = hasLowerBound ? String.valueOf(valueRange.lowerEndpoint()) : getLowerEndpoint(tableNameCache);
+ String max = hasUpperBound ? String.valueOf(valueRange.upperEndpoint()) : getUpperEndpoint(tableNameCache);
+ // 循环计算分表范围
+ Set resultTableNames = new LinkedHashSet<>();
+ try {
+ Date minDate = DATE_TIME_FORMATTER_SPILE.parse(min);
+ Date maxDate = DATE_TIME_FORMATTER_SPILE.parse(max);
+ Calendar calendar = Calendar.getInstance();
+ while (minDate.before(maxDate) || minDate.equals(maxDate)) {
+ String tableName = logicTableName + TABLE_SPLIT_SYMBOL + TABLE_SHARD_Date_FORMATTER.format(minDate);
+ resultTableNames.add(tableName);
+ calendar.setTime(minDate); // 设置Calendar的时间为Date对象的时间
+ calendar.add(Calendar.DAY_OF_MONTH, 1); // 给日期加一天
+ minDate = calendar.getTime();
+ }
+ log.info("开始分表查询结束:{}",System.currentTimeMillis());
+ return ShardingAlgorithmTool.getShardingTablesAndCreate(logicTable, resultTableNames);
+ } catch (Exception e) {
+ return ShardingAlgorithmTool.getShardingTablesAndCreate(logicTable, logicTable.resultTableNamesCache());
+ }
+ }
+
+
+ @Override
+ public void init() {
+
+ }
+
+ @Override
+ public String getType() {
+ return null;
+ }
+
+ // --------------------------------------------------------------------------------------------------------------
+ // 私有方法
+ // --------------------------------------------------------------------------------------------------------------
+
+ /**
+ * 获取 最小分片值
+ * @param tableNames 表名集合
+ * @return 最小分片值
+ */
+ private String getLowerEndpoint(Collection tableNames) {
+ Optional optional = tableNames.stream()
+ .map(o -> LocalDateTime.parse(o.replace(TABLE_SPLIT_SYMBOL, "") + "01 00:00:00", DATE_TIME_FORMATTER))
+ .min(Comparator.comparing(Function.identity()));
+ if (optional.isPresent()) {
+ ZonedDateTime zonedDateTime = optional.get().atZone(ZoneId.systemDefault());
+ Instant instant = zonedDateTime.toInstant();
+ return String.valueOf(Date.from(instant));
+ } else {
+ log.error(">>>>>>>>>> 【ERROR】获取数据最小分表失败,请稍后重试,tableName:{}", tableNames);
+ throw new IllegalArgumentException("获取数据最小分表失败,请稍后重试");
+ }
+ }
+
+ /**
+ * 获取 最大分片值
+ * @param tableNames 表名集合
+ * @return 最大分片值
+ */
+ private String getUpperEndpoint(Collection tableNames) {
+ Optional optional = tableNames.stream()
+ .map(o -> LocalDateTime.parse(o.replace(TABLE_SPLIT_SYMBOL, "") + "01 00:00:00", DATE_TIME_FORMATTER))
+ .max(Comparator.comparing(Function.identity()));
+ if (optional.isPresent()) {
+ ZonedDateTime zonedDateTime = optional.get().atZone(ZoneId.systemDefault());
+ Instant instant = zonedDateTime.toInstant();
+ return String.valueOf(Date.from(instant));
+ } else {
+ log.error(">>>>>>>>>> 【ERROR】获取数据最大分表失败,请稍后重试,tableName:{}", tableNames);
+ throw new IllegalArgumentException("获取数据最大分表失败,请稍后重试");
+ }
+ }
+
+ /**
+ * 根据分片规则获取的表,创建所有的表
+ * @param logicTable
+ * @param tableNames
+ */
+ private void createAllTable(ShardingTableCacheEnum logicTable, Collection tableNames) {
+ if (!CollectionUtils.isEmpty(logicTable.resultTableNamesCache())) {
+ //如果缓存中有表了,则证明已经创建了表,无需再创建
+ return;
+ }
+ //根据分片规则创建表
+ ShardingAlgorithmTool.getShardingTablesAndCreate(logicTable,tableNames);
+ //刷新缓存
+ ShardingAlgorithmTool.tableNameCacheReload(logicTable);
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/config/sharding/enums/ShardingTableCacheEnum.java b/maibu-framework/src/main/java/com/maibu/config/sharding/enums/ShardingTableCacheEnum.java
new file mode 100644
index 0000000..c47fc75
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/config/sharding/enums/ShardingTableCacheEnum.java
@@ -0,0 +1,83 @@
+package com.maibu.config.sharding.enums;
+
+import com.baomidou.mybatisplus.core.toolkit.CollectionUtils;
+import java.util.*;
+
+import static com.fastbee.framework.config.sharding.ShardingAlgorithmTool.actualDataNodesRefresh;
+
+
+/**
+ * @Title ShardingTableCacheEnum
+ *
@Description 分片表缓存枚举
+ *
+ */
+public enum ShardingTableCacheEnum {
+
+ /**
+ * 用户埋点表
+ */
+ DEVICE_LOG("iot_device_log", new HashSet<>());
+
+ /**
+ * 逻辑表名
+ */
+ private final String logicTableName;
+ /**
+ * 实际表名
+ */
+ private final Set resultTableNamesCache;
+
+ private static Map valueMap = new HashMap<>();
+
+ static {
+ Arrays.stream(ShardingTableCacheEnum.values()).forEach(o -> valueMap.put(o.logicTableName, o));
+ }
+
+ ShardingTableCacheEnum(String logicTableName, Set resultTableNamesCache) {
+ this.logicTableName = logicTableName;
+ this.resultTableNamesCache = resultTableNamesCache;
+ }
+
+ public static ShardingTableCacheEnum of(String value) {
+ return valueMap.get(value);
+ }
+
+ public String logicTableName() {
+ return logicTableName;
+ }
+
+ public Set resultTableNamesCache() {
+ return resultTableNamesCache;
+ }
+
+ /**
+ * 更新缓存、配置(原子操作)
+ *
+ * @param tableNameList
+ */
+ public void atomicUpdateCacheAndActualDataNodes(List tableNameList) {
+ if (CollectionUtils.isEmpty(tableNameList)) {
+ return;
+ }
+ synchronized (resultTableNamesCache) {
+ // 删除缓存
+ resultTableNamesCache.clear();
+ // 写入新的缓存
+ resultTableNamesCache.addAll(tableNameList);
+ // 动态更新配置 actualDataNodes
+ actualDataNodesRefresh(logicTableName, tableNameList);
+ }
+ }
+
+ public static Set logicTableNames() {
+ return valueMap.keySet();
+ }
+
+ @Override
+ public String toString() {
+ return "ShardingTableCacheEnum{" +
+ "logicTableName='" + logicTableName + '\'' +
+ ", resultTableNamesCache=" + resultTableNamesCache +
+ '}';
+ }
+}
diff --git a/maibu-framework/src/main/java/com/maibu/datasource/DynamicDataSource.java b/maibu-framework/src/main/java/com/maibu/datasource/DynamicDataSource.java
new file mode 100644
index 0000000..ae408f3
--- /dev/null
+++ b/maibu-framework/src/main/java/com/maibu/datasource/DynamicDataSource.java
@@ -0,0 +1,30 @@
+package com.maibu.datasource;
+
+import java.util.Map;
+import javax.sql.DataSource;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.jdbc.datasource.lookup.AbstractRoutingDataSource;
+
+/**
+ * 动态数据源
+ *
+ * @author ruoyi
+ */
+@Slf4j
+public class DynamicDataSource extends AbstractRoutingDataSource
+{
+ public DynamicDataSource(DataSource defaultTargetDataSource, Map