咨询:并发处理时如何将数据刷入数据库?附多线程入队代码
针对你现在8个线程并发处理数据后刷入数据库的场景,我结合你的代码结构给你几个实用的落地方案:
方案1:单独启动数据库消费线程(推荐)
这种方式符合生产者-消费者模型,把业务数据生产和数据库IO操作解耦,避免业务线程被IO阻塞,也更容易控制入库的批量和节奏。
改造步骤:
替换共享集合为线程安全阻塞队列
把原来的tagProcess.getAutotagResult()(如果是普通List的话)换成ArrayBlockingQueue<String[]>,让8个PairProcessor线程作为生产者往队列里塞数据。实现独立的数据库写入线程
写一个专门的DbWriter线程,循环从队列取数据,攒到指定批量后执行批量插入:
public class DbWriter implements Runnable { private BlockingQueue<String[]> dataQueue; private JdbcTemplate jdbcTemplate; // 用Spring JdbcTemplate或自定义连接池 private static final int BATCH_SIZE = 100; // 可根据数据库性能调整 public DbWriter(BlockingQueue<String[]> dataQueue, JdbcTemplate jdbcTemplate) { this.dataQueue = dataQueue; this.jdbcTemplate = jdbcTemplate; } @Override public void run() { List<String[]> batchData = new ArrayList<>(BATCH_SIZE); while (true) { try { // 阻塞等待数据,也可以用poll加超时判断是否停止 String[] row = dataQueue.take(); if (row == null) { // 用null作为生产结束的标志 // 先把剩余数据入库再退出 if (!batchData.isEmpty()) { insertBatch(batchData); } break; } batchData.add(row); // 达到批量阈值就执行入库 if (batchData.size() >= BATCH_SIZE) { insertBatch(batchData); batchData.clear(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 中断时先处理剩余数据 if (!batchData.isEmpty()) { insertBatch(batchData); } break; } } } private void insertBatch(List<String[]> batchData) { String sql = "INSERT INTO your_table(col1, col2, ...) VALUES (?, ?, ...)"; jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int i) throws SQLException { String[] row = batchData.get(i); ps.setString(1, row[0]); ps.setString(2, row[1]); // 按你的表字段设置对应参数 } @Override public int getBatchSize() { return batchData.size(); } }); } }
- 修改PairProducer代码
不再往共享List加数据,直接往阻塞队列推送:
@Override public void run() { List<String[]> twoRow; while (true) { twoRow = tagProcess.getNext(); if (twoRow == null) break; try { for (String[] row : twoRow) { dataQueue.put(row); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } // 生产结束后往队列放null(如果有8个生产者,需要放8个null,或者用CountDownLatch统一通知) try { dataQueue.put(null); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }
关键注意点:
- 用
CountDownLatch等待所有生产者线程完成,确保消费者能收到所有结束信号 - 数据库必须用连接池(比如HikariCP),绝对不能每次插入新建连接
- 批量大小建议在50-200之间,根据你的数据库写入性能调整
方案2:每个线程本地攒批量,达到阈值后入库
如果不想单独开消费线程,可以让每个PairProcessor自己在本地攒数据,到指定数量后执行批量插入,避免共享集合的线程安全问题。
改造后的PairProcessor代码:
public class PairProcessor implements Runnable { private TagProcess tagProcess; private JdbcTemplate jdbcTemplate; private static final int BATCH_SIZE = 100; // 每个线程用自己的本地List,避免线程安全问题 private List<String[]> localBatch = new ArrayList<>(BATCH_SIZE); public PairProcessor(TagProcess tagProcess, JdbcTemplate jdbcTemplate) { this.tagProcess = tagProcess; this.jdbcTemplate = jdbcTemplate; } @Override public void run() { List<String[]> twoRow; while (true) { twoRow = tagProcess.getNext(); if (twoRow == null) break; localBatch.addAll(twoRow); // 达到批量阈值就入库 if (localBatch.size() >= BATCH_SIZE) { insertBatch(localBatch); localBatch.clear(); } } // 线程结束前把剩余数据入库 if (!localBatch.isEmpty()) { insertBatch(localBatch); } } private void insertBatch(List<String[]> batchData) { // 和方案1的批量插入代码一致 String sql = "INSERT INTO your_table(col1, col2, ...) VALUES (?, ?, ...)"; jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int i) throws SQLException { String[] row = batchData.get(i); ps.setString(1, row[0]); ps.setString(2, row[1]); } @Override public int getBatchSize() { return batchData.size(); } }); } }
关键注意点:
- 确保数据库连接池的最大连接数足够(建议设置为线程数+5)
- 如果数据库有写入瓶颈,这种多线程同时写的方式可能不如单消费线程可控
通用核心注意事项
- 事务控制:如果要求批量数据要么全成功要么全失败,给批量插入添加事务(比如Spring的
@Transactional,或手动管理Connection事务) - 异常重试:入库失败时要加重试机制(比如用Guava Retryer),或者把失败数据放到死信队列后续处理,不能直接丢弃
- 资源释放:所有线程结束后,务必关闭数据库连接池,避免资源泄漏
内容的提问来源于stack exchange,提问作者user4079032
相关产品推荐
相关产品推荐

