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

咨询:并发处理时如何将数据刷入数据库?附多线程入队代码

针对你现在8个线程并发处理数据后刷入数据库的场景,我结合你的代码结构给你几个实用的落地方案:

方案1:单独启动数据库消费线程(推荐)

这种方式符合生产者-消费者模型,把业务数据生产和数据库IO操作解耦,避免业务线程被IO阻塞,也更容易控制入库的批量和节奏。

改造步骤:

  1. 替换共享集合为线程安全阻塞队列
    把原来的tagProcess.getAutotagResult()(如果是普通List的话)换成ArrayBlockingQueue<String[]>,让8个PairProcessor线程作为生产者往队列里塞数据。

  2. 实现独立的数据库写入线程
    写一个专门的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();
            }
        });
    }
}
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:24:21