Skip to content

ES数据批量导入

一、核心工具介绍

1. CountDownLatch 倒计时闭锁

作用

实现线程同步等待:主线程等待多个子线程全部完成任务后,再继续执行后续逻辑

核心API

  1. CountDownLatch(int count):构造方法,传入需要等待的子线程总数,初始化计数器;
  2. await():主线程阻塞等待,直到计数器归0才会解除阻塞;
  3. countDown():子线程执行完毕后调用,计数器数值-1。

执行流程

主线程调用await()进入阻塞;多个业务子线程执行完任务后分别执行countDown();每调用一次计数器减1,全部减到0后主线程唤醒继续运行。

2. Future

搭配线程池submit()方法使用,用于接收异步任务返回结果,支持阻塞等待任务完成、获取执行返回值。

二、实战业务场景:千万级DB数据批量导入ES

业务背景

项目上线初始化,需要将数据库约1000万条全量数据同步至Elasticsearch索引;一次性全量读取会加载大量数据到内存,直接触发OOM内存溢出,因此采用分页分片+线程池并发方案,搭配CountDownLatch控制全量同步完成。

完整业务执行流程

  1. 查询DB总数据条数,设定单页分页容量(示例每页2000条),计算总分页数;
  2. 根据总分页数初始化CountDownLatch,计数器=总页数;
  3. 循环分页查询单页DB数据,每一页封装为独立导入任务提交至自定义线程池;
  4. 线程池内子线程执行ES批量写入,写入完成调用countDown()将计数器-1;
  5. 主线程调用await()阻塞,等待所有分页导入任务全部执行完毕(计数器归0);
  6. 主线程解除阻塞,执行同步完成后的收尾逻辑(如打印同步耗时、更新同步状态标记)。

方案优势

  1. 分页分片读取,避免一次性加载千万数据,防止OOM;
  2. 线程池并发写入ES,相比单线程同步大幅缩短全量同步耗时;
  3. CountDownLatch精准控制主线程等待全部异步任务完成,保证同步流程闭环。

三、其他常见项目多线程场景

  1. 多接口并行查询(Future) 页面需要同时调用多个互不依赖的第三方接口/数据库查询,使用线程池submit()并行执行,通过Future批量获取返回数据,降低接口总响应时间。
  2. 异步消息发送 订单创建、支付完成后,异步发送短信、站内信、MQ消息,不阻塞主业务流程。
  3. 文件批量处理 批量解析/上传/导出大量文件,多线程分担IO密集操作,提升处理速度。
  4. 定时任务分片执行 大数据定时统计任务,拆分分片交由线程池并发计算。

四、面试标准回答(项目哪里用到多线程)

在项目初始化全量同步数据库千万级数据到ES的场景中使用了多线程:

  1. 需求:一次性读取全部数据会OOM,需要分页并发导入提升同步速度;
  2. 技术选型:手动创建自定义参数线程池,搭配CountDownLatch做线程同步;
  3. 实现逻辑:分页拆分数据,每页封装导入任务提交线程池,子线程导入完成执行countDown,主线程await等待全部分页任务结束后再执行收尾;
  4. 收益:解决全量加载内存溢出问题,并发写入缩短数据同步耗时,同步流程可控。

Powered by VitePress 1.6.4 | 持续更新中