MyBatis-Plus大数据量同步:分页+BATCH模式
·
该方案事务是全部分页执行完提交,整个同步过程是一个大事务,要么全部成功,要么全部回滚。适合五万条以下数据量。
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.example.domain.SyncSource;
import com.example.domain.SyncTarget;
import com.example.domain.result.AjaxResult;
import com.example.mapper.SyncSourceMapper;
import com.example.mapper.SyncTargetMapper;
import com.example.service.ISyncTargetService;
import org.apache.ibatis.session.ExecutorType;
import org.apache.ibatis.session.SqlSession;
import org.apache.ibatis.session.SqlSessionFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.stream.Collectors;
@Service
public class SyncTargetServiceImpl implements ISyncTargetService {
private static final Logger log = LoggerFactory.getLogger(SyncTargetServiceImpl.class);
@Autowired
private SqlSessionFactory sqlSessionFactory;
@Override
public AjaxResult sync() {
long startTime = System.currentTimeMillis();
log.info("========== 开始数据同步 ==========");
// 配置参数
int pageSize = 500;
int pageNum = 1;
int totalInsert = 0;
int totalUpdate = 0;
int totalPage = 0;
// 使用BATCH执行器,包住所有操作
try (SqlSession sqlSession = sqlSessionFactory.openSession(ExecutorType.BATCH)) {
// 从同一个SqlSession获取所有Mapper
SyncSourceMapper sourceMapper = sqlSession.getMapper(SyncSourceMapper.class);
SyncTargetMapper targetMapper = sqlSession.getMapper(SyncTargetMapper.class);
while (true) {
log.info("---------- 正在处理第{}页数据 ----------", pageNum);
// 使用BATCH的Mapper查询源数据
int offset = (pageNum - 1) * pageSize;
List<SyncSource> sourceList = sourceMapper.selectList(
new LambdaQueryWrapper<SyncSource>()
.orderByAsc(SyncSource::getId)
.last("LIMIT " + offset + ", " + pageSize)
);
if (sourceList.isEmpty()) {
log.info("没有更多数据,同步完成");
break;
}
totalPage++;
log.info("第{}页查询到{}条数据", pageNum, sourceList.size());
// 使用BATCH的Mapper查询已存在的ID
Set<Long> batchIds = sourceList.stream()
.map(SyncSource::getId)
.collect(Collectors.toSet());
List<SyncTarget> existingList = targetMapper.selectList(
new LambdaQueryWrapper<SyncTarget>()
.select(SyncTarget::getId)
.in(SyncTarget::getId, batchIds)
);
Set<Long> existingIds = existingList.stream()
.map(SyncTarget::getId)
.collect(Collectors.toSet());
log.info("当前批次{}个ID中,目标表已存在{}个", batchIds.size(), existingIds.size());
// 分流:区分新增和更新
List<SyncTarget> insertList = new ArrayList<>();
List<SyncTarget> updateList = new ArrayList<>();
for (SyncSource source : sourceList) {
SyncTarget target = convertToTarget(source);
if (existingIds.contains(target.getId())) {
// 已存在 → 更新
target.setUpdateTime(LocalDateTime.now());
updateList.add(target);
} else {
// 不存在 → 新增
target.setCreateTime(LocalDateTime.now());
target.setUpdateTime(LocalDateTime.now());
insertList.add(target);
}
}
log.info("第{}页分流结果:新增{}条,更新{}条", pageNum, insertList.size(), updateList.size());
// 直接使用BATCH的Mapper
if (!insertList.isEmpty()) {
for (SyncTarget target : insertList) {
targetMapper.insert(target);
}
}
if (!updateList.isEmpty()) {
for (SyncTarget target : updateList) {
targetMapper.updateById(target);
}
}
// 每页处理完后刷新(但不提交)
sqlSession.flushStatements();
totalInsert += insertList.size();
totalUpdate += updateList.size();
pageNum++;
if (sourceList.size() < pageSize) {
break;
}
}
// 所有数据处理完成,统一提交
sqlSession.commit();
log.info("所有数据提交成功");
long endTime = System.currentTimeMillis();
log.info("========== 数据同步完成!总耗时:{}ms ==========", (endTime - startTime));
log.info("========== 共处理{}页数据 ==========", totalPage);
log.info("========== 新增{}条,更新{}条 ==========", totalInsert, totalUpdate);
return AjaxResult.success();
} catch (Exception e) {
log.error("数据同步失败,事务将自动回滚", e);
return AjaxResult.error("同步失败:" + e.getMessage());
}
}
/**
* 将源数据转换为目标数据
*/
private SyncTarget convertToTarget(SyncSource source) {
SyncTarget target = new SyncTarget();
target.setId(source.getId());
// 如果有其他字段,继续映射
return target;
}
}
更多推荐

所有评论(0)