一种轻量分表方案-MyBatis拦截器分表实践

背景

部门内有一些亿级别核心业务表增速非常快,增量日均100W,但线上业务只依赖近一周的数据。随着数据量的迅速增长,慢SQL频发,数据库性能下降,系统稳定性受到严重影响。本篇文章,将分享如何使用MyBatis拦截器低成本的提升数据库稳定性。

业界常见方案

针对冷数据多的大表,常用的策略有以2种:

1. 删除/归档旧数据。

2. 分表。

归档/删除旧数据

定期将冷数据移动到归档表或者冷存储中,或定期对表进行删除,以减少表的大小。此策略逻辑简单,只需要编写一个JOB定期执行SQL删除数据。我们开始也是用这种方案,但此方案也有一些副作用:

1.数据删除会影响数据库性能,引发慢sql,多张表并行删除,数据库压力会更大。

2.频繁删除数据,会产生数据库碎片,影响数据库性能,引发慢SQL。

综上,此方案有一定风险,为了规避这种风险,我们决定采用另一种方案:分表。

分表

我们决定按日期对表进行横向拆分,实现让系统每周生成一张周期表,表内只存近一周的数据,规避单表过大带来的风险。

分表方案选型

经调研,考虑2种分表方案:Sharding-JDBC、利用Mybatis自带的拦截器特性。

经过对比后,决定采用Mybatis拦截器来实现分表,原因如下:

1.JAVA生态中很常用的分表框架是Sharding-JDBC,虽然功能强大,但需要一定的接入成本,并且很多功能暂时用不上。

2.系统本身已经在使用Mybatis了,只需要添加一个mybaits拦截器,把SQL表名替换为新的周期表就可以了,没有接入新框架的成本,开发成本也不高。

简易架构图

分表具体实现代码

分表配置对象

1import lombok.AllArgsConstructor; 2import lombok.Data; 3import lombok.NoArgsConstructor; 4 5import java.util.Date; 6 7@Data 8@AllArgsConstructor 9@NoArgsConstructor 10public class ShardingProperty { 11 // 分表周期天数,配置7,就是一周一分 12 private Integer days; 13 // 分表开始日期,需要用这个日期计算周期表名 14 private Date beginDate; 15 // 需要分表的表名 16 private String tableName; 17} 18 19 20

分表配置类

1import java.util.concurrent.ConcurrentHashMap; 2 3public class ShardingPropertyConfig { 4 5 public static final ConcurrentHashMap<String, ShardingProperty> SHARDING_TABLE = new ConcurrentHashMap<>(); 6 7 static { 8 ShardingProperty orderInfoShardingConfig = new ShardingProperty(15, DateUtils.string2Date("20231117"), "order_info"); 9 ShardingProperty userInfoShardingConfig = new ShardingProperty(7, DateUtils.string2Date("20231117"), "user_info"); 10 11 SHARDING_TABLE.put(orderInfoShardingConfig.getTableName(), orderInfoShardingConfig); 12 SHARDING_TABLE.put(userInfoShardingConfig.getTableName(), userInfoShardingConfig); 13 } 14} 15 16

拦截器

1import lombok.extern.slf4j.Slf4j; 2import o2o.aspect.platform.function.template.service.TemplateMatchService; 3import org.apache.commons.lang3.StringUtils; 4import org.apache.ibatis.executor.statement.StatementHandler; 5import org.apache.ibatis.mapping.BoundSql; 6import org.apache.ibatis.mapping.MappedStatement; 7import org.apache.ibatis.plugin.*; 8import org.apache.ibatis.reflection.DefaultReflectorFactory; 9import org.apache.ibatis.reflection.MetaObject; 10import org.apache.ibatis.reflection.ReflectorFactory; 11import org.apache.ibatis.reflection.factory.DefaultObjectFactory; 12import org.apache.ibatis.reflection.factory.ObjectFactory; 13import org.apache.ibatis.reflection.wrapper.DefaultObjectWrapperFactory; 14import org.apache.ibatis.reflection.wrapper.ObjectWrapperFactory; 15import org.springframework.stereotype.Component; 16 17import java.sql.Connection; 18import java.time.LocalDateTime; 19import java.time.format.DateTimeFormatter; 20import java.util.Date; 21import java.util.Properties; 22 23@Slf4j 24@Component 25@Intercepts({@Signature(type = StatementHandler.class, method = "prepare", args = {Connection.class, Integer.class})}) 26public class ShardingTableInterceptor implements Interceptor { 27 private static final ObjectFactory DEFAULT_OBJECT_FACTORY = new DefaultObjectFactory(); 28 private static final ObjectWrapperFactory DEFAULT_OBJECT_WRAPPER_FACTORY = new DefaultObjectWrapperFactory(); 29 private static final ReflectorFactory DEFAULT_REFLECTOR_FACTORY = new DefaultReflectorFactory(); 30 private static final String MAPPED_STATEMENT = "delegate.mappedStatement"; 31 private static final String BOUND_SQL = "delegate.boundSql"; 32 private static final String ORIGIN_BOUND_SQL = "delegate.boundSql.sql"; 33 private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyyMMdd"); 34 private static final String SHARDING_MAPPER = "com.jd.o2o.inviter.promote.mapper.ShardingMapper"; 35 36 private ConfigUtils configUtils = SpringContextHolder.getBean(ConfigUtils.class); 37 38 @Override 39 public Object intercept(Invocation invocation) throws Throwable { 40 boolean shardingSwitch = configUtils.getBool("sharding_switch", false); 41 // 没开启分表 直接返回老数据 42 if (!shardingSwitch) { 43 return invocation.proceed(); 44 } 45 46 StatementHandler statementHandler = (StatementHandler) invocation.getTarget(); 47 MetaObject metaStatementHandler = MetaObject.forObject(statementHandler, DEFAULT_OBJECT_FACTORY, DEFAULT_OBJECT_WRAPPER_FACTORY, DEFAULT_REFLECTOR_FACTORY); 48 MappedStatement mappedStatement = (MappedStatement) metaStatementHandler.getValue(MAPPED_STATEMENT); 49 BoundSql boundSql = (BoundSql) metaStatementHandler.getValue(BOUND_SQL); 50 String originSql = (String) metaStatementHandler.getValue(ORIGIN_BOUND_SQL); 51 if (StringUtils.isBlank(originSql)) { 52 return invocation.proceed(); 53 } 54 55 // 获取表名 56 String tableName = TemplateMatchService.matchTableName(boundSql.getSql().trim()); 57 ShardingProperty shardingProperty = ShardingPropertyConfig.SHARDING_TABLE.get(tableName); 58 if (shardingProperty == null) { 59 return invocation.proceed(); 60 } 61 62 // 新表 63 String shardingTable = getCurrentShardingTable(shardingProperty, new Date()); 64 String rebuildSql = boundSql.getSql().replace(shardingProperty.getTableName(), shardingTable); 65 metaStatementHandler.setValue(ORIGIN_BOUND_SQL, rebuildSql); 66 if (log.isDebugEnabled()) { 67 log.info("rebuildSQL -> {}", rebuildSql); 68 } 69 70 return invocation.proceed(); 71 } 72 73 @Override 74 public Object plugin(Object target) { 75 if (target instanceof StatementHandler) { 76 return Plugin.wrap(target, this); 77 } 78 return target; 79 } 80 81 @Override 82 public void setProperties(Properties properties) {} 83 84 public static String getCurrentShardingTable(ShardingProperty shardingProperty, Date createTime) { 85 String tableName = shardingProperty.getTableName(); 86 Integer days = shardingProperty.getDays(); 87 Date beginDate = shardingProperty.getBeginDate(); 88 89 Date date; 90 if (createTime == null) { 91 date = new Date(); 92 } else { 93 date = createTime; 94 } 95 if (date.before(beginDate)) { 96 return null; 97 } 98 LocalDateTime targetDate = SimpleDateFormatUtils.convertDateToLocalDateTime(date); 99 LocalDateTime startDate = SimpleDateFormatUtils.convertDateToLocalDateTime(beginDate); 100 LocalDateTime intervalStartDate = DateIntervalChecker.getIntervalStartDate(targetDate, startDate, days); 101 LocalDateTime intervalEndDate = intervalStartDate.plusDays(days - 1); 102 return tableName + "_" + intervalStartDate.format(FORMATTER) + "_" + intervalEndDate.format(FORMATTER); 103 } 104} 105 106

临界点数据不连续问题

分表方案有1个难点需要解决:周期临界点数据不连续。举例:假设要对operate_log(操作日志表)大表进行横向分表,每周一张表,分表明细可看下面表格。

| 第一周(operate_log_20240107_20240108) | 第二周(operate_log_20240108_20240114) | 第三周(operate_log_20240115_20240121) | | 1月1号 ~ 1月7号的数据 | 1月8号 ~ 1月14号的数据 | 1月15号 ~ 1月21号的数据 |

1月8号就是分表临界点,8号需要切换到第二周的表,但8号0点刚切换的时候,表内没有任何数据,这时如果业务需要查近一周的操作日志是查不到的,这样就会引发线上问题。

我决定采用数据冗余的方式来解决这个痛点。每个周期表都冗余一份上个周期的数据,用双倍数据量实现数据滑动的效果,效果见下面表格。

| 第一周(operate_log_20240107_20240108) | 第二周(operate_log_20240108_20240114) | 第三周(operate_log_20240115_20240121) | | 12月25号 ~ 12月31号的数据 | 1月1号 ~ 1月7号的数据 | 1月8号 ~ 1月14号的数据 | | 1月1号 ~ 1月7号的数据 | 1月8号 ~ 1月14号的数据 | 1月15号 ~ 1月21号的数据 |

注:表格内第一行数据就是冗余的上个周期表的数据。

思路有了,接下来就要考虑怎么实现双写(数据冗余到下个周期表),有2种方案:

1.在SQL执行完成返回结果前添加逻辑(可以用AspectJ 或 mybatis拦截器),如果SQL内的表名是当前周期表,就把表名替换为下个周期表,然后再次执行SQL。此方案对业务影响大,相当于串行执行了2次SQL,有性能损耗。

2.监听增量binlog,京东内部有现成的数据订阅中间件DRC,读者也可以使用cannal等开源中间件来代替DRC,原理大同小异,此方案对业务无影响。

方案对比后,选择了对业务性能损耗小的方案二。

监听binlog并双写流程图

监听binlog数据双写注意点

1.提前上线监听程序,提前把老表数据同步到新的周期表。分表前只监听老表binlog就可以,分表前只需要把老表数据同步到新表。

2.切换到新表的临界点,为了避免丢失积压的老表binlog,需要同时处理新表binlog和老表binlog,这样会出现死循环同步的问题,因为老表需要同步新表,新表又需要双写老表。为了打破循环,需要先把双写老表消费堵上让消息暂时积压,切换新表成功后,再打开双写消费。

监听binlog数据双写代码

注:下面代码不能直接用,只提供基本思路

1/** 2 * 监听binlog ,分表双写,解决数据临界问题 3*/ 4@Slf4j 5@Component 6public class BinLogConsumer implements MessageListener { 7 8 private MessageDeserialize deserialize = new JMQMessageDeserialize(); 9 10 private static final String TABLE_PLACEHOLDER = "%TABLE%"; 11 12 @Value("${mq.doubleWriteTopic.topic}") 13 private String doubleWriteTopic; 14 15 @Autowired 16 private JmqProducerService jmqProducerService; 17 18 19 @Override 20 public void onMessage(List<Message> messages) throws Exception { 21 if (messages == null || messages.isEmpty()) { 22 return; 23 } 24 List<EntryMessage> entryMessages = deserialize.deserialize(messages); 25 for (EntryMessage entryMessage : entryMessages) { 26 try { 27 syncData(entryMessage); 28 } catch (Exception e) { 29 log.error("sharding sync data error", e); 30 throw e; 31 } 32 } 33 } 34 35 private void syncData(EntryMessage entryMessage) throws JMQException { 36 // 根据binlog内的表名,获取需要同步的表 37 // 3种情况: 38 // 1、老表:需要同步当前周期表,和下个周期表。 39 // 2、当前周期表:需要同步下个周期表,和老表。 40 // 3、下个周期表:不需要同步。 41 List<String> syncTables = getSyncTables(entryMessage.tableName, entryMessage.createTime); 42 43 if (CollectionUtils.isEmpty(syncTables)) { 44 log.info("table {} is not need sync", tableName); 45 return; 46 } 47 48 if (entryMessage.getHeader().getEventType() == WaveEntry.EventType.INSERT) { 49 String insertTableSqlTemplate = parseSqlForInsert(rowData); 50 for (String syncTable : syncTables) { 51 String insertSql = insertTableSqlTemplate.replaceAll(TABLE_PLACEHOLDER, syncTable); 52 // 双写老表发Q,为了避免出现同步死循环问题 53 if (ShardingPropertyConfig.SHARDING_TABLE.containsKey(syncTable)) { 54 Long primaryKey = getPrimaryKey(rowData.getAfterColumnsList()); 55 sendDoubleWriteMsg(insertSql, primaryKey); 56 continue; 57 } 58 mysqlConnection.executeSql(insertSql); 59 } 60 continue; 61 } 62 } 63 64 65

数据对比

为了保证新表和老表数据一致,需要编写对比程序,在上线前进行数据对比,保证binlog同步无问题。

具体实现代码不做展示,思路:新表查询一定量级数据,老表查询相同量级数据,都转换成JSON,equals对比。

作者:京东零售 张均杰

来源:京东云开发者社区 转载请注明来源

点赞
收藏

评论区

加载中...

相关推荐

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

sql注入

反引号是个比较特别的字符,下面记录下怎么利用0x00SQL注入反引号可利用在分隔符及注释作用,不过使用范围只于表名、数据库名、字段名、起别名这些场景,下面具体说下1)表名payload:select\from\users\whereuser\_id1limit0,1;!(https://o

MySQL大表优化方案

背景阿里云RDSFORMySQL(MySQL5.7版本)数据库业务表每月新增数据量超过千万,随着数据量持续增加,我们业务出现大表慢查询,在业务高峰期主业务表的慢查询需要几十秒严重影响业务方案概述!(https://oscimg.oschina.net/oscnet/025a5adc414b403cbb5d3f60b

MySQL 快速创建千万级测试数据

备注:此文章的数据量在100W,如果想要千万级,调大数量即可,但是不要大量使用rand()或者uuid()会导致性能下降背景在进行查询操作的性能测试或者sql优化时,我们经常需要在线下环境构建大量的基础数据供我们测试,模拟线上的真实环境。废话,总不能让我去线上去测试吧,会被DBA砍死的创建测试数据的方式

分库分表后复杂查询的应对之道:基于DTS实时性ES宽表构建技术实践

作者:京东物流王军1问题域业务发展的初期,我们的数据库架构往往是单库单表,外加读写分离来快速的支撑业务,随着用户量和订单量的增加,数据库的计算和存储往往会成为我们系统的瓶颈,业界的实践多数采用分而治之的思想:分库分表,通过分库分表应对存系统读写性能瓶颈和存