多线程导致重复ReportStart对象及Java代码性能优化求助
首先你遇到的重复问题核心原因是分片逻辑没有真正生效:你创建32个线程时传入了分片的起始偏移量2000000 * i,但实际查询语句里完全没用到这个参数!所有线程都在查询全局的status=false记录,这就导致多个线程会抢到同一条未处理的数据,重复生成ReportStart对象。
另外,当前的代码在性能上也有不少可以优化的地方,比如单条操作代替批量、事务提交过于频繁等。
方案1:基于ID的分片处理(推荐,无锁竞争)
既然StartCgnatLog的ID是唯一且连续的(从你的分片逻辑来看应该是),让每个线程只处理自己分片范围内的ID,彻底避免线程间的数据重叠:
- 修改任务类,接收分片的起始和结束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() { // 后续处理逻辑用这个分片范围过滤 } }
- 创建线程时计算每个分片的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(); } }
- 修改查询语句,只查询当前分片范围内未处理的记录:
// 预编译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

