1. 注解 @PostConstruct

  • 作用:该注解标记的方法会在 Spring Bean 初始化完成后自动执行。

  • 场景:用于启动后台任务(如定时任务、异步队列消费)。

2. 方法 start()

  • 功能:初始化并启动一个后台线程,用于持续处理队列中的数据。

  • 关键逻辑

new Thread(() -> { 
    // 线程逻辑
}).start();

3. 线程内部逻辑

(1) 日志记录


log.info("document save task started");
• 作用:标记任务已启动。


(2) 初始化列表
java

复制

下载
List<Document> list = new ArrayList<>();
• 用途:临时存储从队列中批量获取的待处理数据。
(3) 循环处理条件 while (isRunning)
• 作用:通过 isRunning 标志控制线程启停。
• 终止条件:当 isRunning 设为 false 时,循环退出,线程结束。
(4) 清空列表
java

复制

下载
list.clear();
• 目的:避免前一次处理的数据残留。
(5) 从队列中阻塞获取数据 queue.take()
java

复制

下载
try {
    data 
= queue.take();  // 阻塞直到队列非空
} catch (InterruptedException e) {
    log
.error("event batch save task error", e);
}
• 队列类型:queue 应为 BlockingQueue<Document>(如 LinkedBlockingQueue)。
• 数据来源:其他线程通过 queue.put(document) 或 queue.offer(document) 向队列中添加数据。
示例场景:
java

复制

下载
// 其他代码(如 Controller 或 Service)
public void addDocumentToQueue(Document doc) {
    queue
.put(doc);  // 数据被添加到队列
}
(6) 批量提取数据 queue.drainTo(list, 2000)
java

复制

下载
list.add(data);          // 添加第一个元素
queue
.drainTo(list, 2000);  // 批量提取最多 2000 个元素
• 作用:一次性从队列中取出多个元素(最多 2000 个),减少频繁操作队列的开销。
• 结果:list 包含 1(来自 take()) + 2000(来自 drainTo())个元素。
(7) 调用 refresh(list)
java

复制

下载
refresh(list);  // 批量处理数据(如更新 Redis)
• 用途:实现具体的业务逻辑(如批量持久化到数据库或更新缓存)。

4. queue.take() 数据来源详解
• 生产者-消费者模型:
◦ 生产者:其他代码(如接口请求、定时任务)调用 queue.put(document) 向队列添加数据。
◦ 消费者:当前线程通过 queue.take() 消费数据。
• 队列实现:通常使用线程安全的 BlockingQueue,如:
java

复制

下载
@Autowired
private BlockingQueue<Document> queue;  // 需在 Spring 中初始化

5. 完整流程示意图
图表 代码



消费者(当前线程)
阻塞队列
生产者(其他代码)
消费者(当前线程)
阻塞队列
生产者(其他代码)
loop
[持续处理]
queue.put(document)
queue.take()
返回一个 Document
queue.drainTo(list, 2000)
返回最多 2000 个 Document
refresh(list)

6. 关键注意事项
1. 线程安全:
◦ queue 必须是线程安全的(如 LinkedBlockingQueue)。
◦ isRunning 需声明为 volatile,确保多线程可见性:
java

复制

下载
private volatile boolean isRunning = true;
2. 优雅停机:
在应用关闭时,需设置 isRunning = false 并中断线程:
java

复制

下载
@PreDestroy
public void stop() {
    isRunning 
= false;
    // 可选:中断线程(若线程阻塞在 queue.take())
}
3. 异常处理优化:
◦ 捕获 InterruptedException 后应重新设置中断状态:
java

复制

下载
catch (InterruptedException e) {
    log
.error("Task interrupted", e);
    Thread.currentThread().interrupt();  // 重新设置中断状态
}
4. 性能优化:
◦ 调整 drainTo 的批量大小(如 2000)以平衡吞吐量和内存占用。
◦ 使用线程池(如 ExecutorService)替代直接创建 Thread。

7. 总结
• queue.take() 从 BlockingQueue 中阻塞获取数据,数据由其他线程通过 queue.put() 添加。
• 代码核心逻辑:通过生产者-消费者模型实现异步批量处理(如文档保存)。
• 改进方向:优化线程管理、增加优雅停机逻辑、完善异常处理。

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区