该方案事务是全部分页执行完提交,整个同步过程是一个大事务,要么全部成功,要么全部回滚。适合五万条以下数据量。

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;
    }
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐