You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多线程导致重复ReportStart对象及Java代码性能优化求助

问题分析

首先你遇到的重复问题核心原因是分片逻辑没有真正生效:你创建32个线程时传入了分片的起始偏移量2000000 * i,但实际查询语句里完全没用到这个参数!所有线程都在查询全局的status=false记录,这就导致多个线程会抢到同一条未处理的数据,重复生成ReportStart对象。

另外,当前的代码在性能上也有不少可以优化的地方,比如单条操作代替批量、事务提交过于频繁等。


解决重复处理问题

方案1:基于ID的分片处理(推荐,无锁竞争)

既然StartCgnatLog的ID是唯一且连续的(从你的分片逻辑来看应该是),让每个线程只处理自己分片范围内的ID,彻底避免线程间的数据重叠:

  1. 修改任务类,接收分片的起始和结束ID:
public class ReportStartCreatorTask implements Runnable {
    private Session session;
    private long startId;
    private long endId;

    public ReportStartCreatorTask(Session session, long startId, long endId) {
        this.session = session;
        this.startId = startId;
        this.endId = endId;
    }

    @Override
    public void run() {
        // 后续处理逻辑用这个分片范围过滤
    }
}
  1. 创建线程时计算每个分片的ID范围:
private void createReportStartTask() { 
    SessionFactory sessionFactory = HibernateUtil.getSessionFactory(); 
    ExecutorService service = Executors.newFixedThreadPool(32); 
    long batchSize = 2000000L;
    for (int i = 0; i < 32; i++) { 
        long startId = batchSize * i;
        long endId = batchSize * (i + 1) - 1;
        service.submit(new ReportStartCreatorTask(sessionFactory.openSession(), startId, endId)); 
    }
    // 记得关闭线程池并等待任务完成
    service.shutdown();
    try {
        service.awaitTermination(1, TimeUnit.HOURS);
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
    }
}
  1. 修改查询语句,只查询当前分片范围内未处理的记录:
// 预编译Query,提升性能
Query<StartCgnatLog> query = session.createQuery(
    "from StartCgnatLog where id between :startId and :endId and status = false",
    StartCgnatLog.class
).setMaxResults(500);

session.beginTransaction();
List<StartCgnatLog> startCgnatLogList = query.setParameter("startId", startId)
                                             .setParameter("endId", endId)
                                             .getResultList();

while (!startCgnatLogList.isEmpty()) { 
    // 处理逻辑...
    // 重新查询时复用同一个Query
    startCgnatLogList = query.setParameter("startId", startId)
                             .setParameter("endId", endId)
                             .getResultList();
}

方案2:行级锁(适合ID不连续的场景)

如果ID不是连续的,无法用分片,可以在查询时使用行级锁,避免多个线程拿到同一条数据。以MySQL为例,使用FOR UPDATE SKIP LOCKED(MySQL 8.0+支持):

Query<StartCgnatLog> query = session.createQuery(
    "from StartCgnatLog where status = false for update skip locked",
    StartCgnatLog.class
).setMaxResults(500);

这样每个线程只会拿到未被其他线程锁定的记录,不会重复处理。但这种方式会有锁竞争,性能不如分片方案。


性能优化建议

1. 批量操作替代单条操作

当前代码是循环单条save和update,效率极低。改用Hibernate的批量操作:

// 批量保存ReportStart
List<ReportStart> reportStartList = new ArrayList<>(500);
// 批量更新StartCgnatLog
List<StartCgnatLog> updateList = new ArrayList<>(500);

for (StartCgnatLog log : startCgnatLogList) {
    ReportStart reportStart = new ReportStart();
    reportStart.setStartId(log.getId());
    reportStartList.add(reportStart);
    
    log.setStatus(true);
    updateList.add(log);
}

// 批量插入
session.saveAll(reportStartList);
// 批量更新
session.saveAll(updateList); // 或者使用updateAll,取决于Hibernate版本

同时在Hibernate配置中开启批量支持:

hibernate.jdbc.batch_size=500
hibernate.order_inserts=true
hibernate.order_updates=true

2. 优化事务提交频率

当前每处理500条就提交一次事务,频繁的事务提交会带来额外开销。可以每处理N批(比如10批,5000条)再提交一次,但要注意内存占用,避免OOM:

int batchCount = 0;
final int COMMIT_BATCH = 10;

session.beginTransaction();
while (!startCgnatLogList.isEmpty()) { 
    // 批量处理逻辑...
    
    batchCount++;
    if (batchCount >= COMMIT_BATCH) {
        session.getTransaction().commit();
        session.beginTransaction();
        batchCount = 0;
    }
    
    // 清空缓存,避免内存溢出
    session.flush();
    session.clear();
    
    startCgnatLogList = query.getResultList();
}
// 提交剩余的事务
if (batchCount > 0) {
    session.getTransaction().commit();
}

3. 使用高效的连接池

Hibernate自带的连接池性能较差,建议改用HikariCP,配置合适的连接数(一般和线程数相当,比如32):

hibernate.connection.provider_class=com.zaxxer.hikari.hibernate.HikariConnectionProvider
hibernate.hikari.maximumPoolSize=32
hibernate.hikari.minimumIdle=8

4. 索引优化

给status字段添加索引,加速查询:

CREATE INDEX idx_startcgnatlog_status ON start_cgnat_log(status);

如果用分片方案,因为id是主键已经有索引,id between ... and ... and status=false的查询会利用主键索引+组合索引,效率更高。

5. 避免Session泄漏

确保每个线程的Session在任务结束后关闭:

@Override
public void run() {
    try {
        // 处理逻辑...
    } finally {
        if (session.isOpen()) {
            session.close();
        }
    }
}

完整优化后的任务类示例
public class ReportStartCreatorTask implements Runnable {
    private Session session;
    private long startId;
    private long endId;
    private static final int BATCH_SIZE = 500;
    private static final int COMMIT_BATCH = 10;

    public ReportStartCreatorTask(Session session, long startId, long endId) {
        this.session = session;
        this.startId = startId;
        this.endId = endId;
    }

    @Override
    public void run() {
        try {
            Query<StartCgnatLog> query = session.createQuery(
                    "from StartCgnatLog where id between :startId and :endId and status = false",
                    StartCgnatLog.class
            ).setMaxResults(BATCH_SIZE);

            int batchCount = 0;
            session.beginTransaction();

            List<StartCgnatLog> startCgnatLogList = query.setParameter("startId", startId)
                    .setParameter("endId", endId)
                    .getResultList();

            while (!startCgnatLogList.isEmpty()) {
                List<ReportStart> reportStartList = new ArrayList<>(BATCH_SIZE);
                List<StartCgnatLog> updateList = new ArrayList<>(BATCH_SIZE);

                for (StartCgnatLog log : startCgnatLogList) {
                    ReportStart reportStart = new ReportStart();
                    reportStart.setStartId(log.getId());
                    reportStartList.add(reportStart);

                    log.setStatus(true);
                    updateList.add(log);
                }

                // 批量操作
                session.saveAll(reportStartList);
                session.saveAll(updateList);

                batchCount++;
                if (batchCount >= COMMIT_BATCH) {
                    session.getTransaction().commit();
                    session.beginTransaction();
                    batchCount = 0;
                }

                session.flush();
                session.clear();

                // 重新查询
                startCgnatLogList = query.setParameter("startId", startId)
                        .setParameter("endId", endId)
                        .getResultList();
            }

            // 提交剩余事务
            if (batchCount > 0) {
                session.getTransaction().commit();
            }
        } catch (Exception e) {
            if (session.getTransaction().isActive()) {
                session.getTransaction().rollback();
            }
            e.printStackTrace();
        } finally {
            if (session.isOpen()) {
                session.close();
            }
        }
    }
}

内容的提问来源于stack exchange,提问作者yed2393

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 08:45:28