Java中消息隊列任務(wù)的平滑關(guān)閉詳解
前言
消息隊列中間件是分布式系統(tǒng)中重要的組件,主要解決應(yīng)用解耦,異步消息,流量削鋒等問題,實現(xiàn)高性能,高可用,可伸縮和最終一致性架構(gòu)。目前使用較多的消息隊列有ActiveMQ,RabbitMQ,ZeroMQ,Kafka,MetaMQ,RocketMQ
消息隊列應(yīng)用場景
消息隊列在實際應(yīng)用中常用的使用場景:異步處理,應(yīng)用解耦,流量削鋒和消息通訊四個場景。
本文主要給大家介紹的是關(guān)于Java中消息隊列任務(wù)平滑關(guān)閉的相關(guān)內(nèi)容,分享出來供大家參考學習,下面話不多說了,來一起看看詳細的介紹吧。
1.問題背景
對于消息隊列任務(wù)的監(jiān)聽,我們一般使用Java寫一個獨立的程序,在Linux服務(wù)器上運行。當訂閱者程序啟動后,會通過消息隊列客戶端接收消息,放入線程池中并發(fā)的處理。
那么問題來了,當我們修改程序后,需要重新啟動時,如何保證消息都能夠被處理呢?
一些開源的消息隊列中間件,會提供ACK機制(消息確認機制),當訂閱者處理完消息后,會通知服務(wù)端刪除對應(yīng)消息,如果訂閱者出現(xiàn)異常,服務(wù)端未收到確認消費,則會重試發(fā)送。
那如果消息隊列中間件沒有提供ACK機制,或者為了高吞度量的考慮關(guān)閉了ACK功能,如何最大可能保證消息都能夠被處理呢?
正常來說,訂閱者程序關(guān)閉后,消息會在隊列中堆積,等待訂閱者下次訂閱消費,所以未接收的消息是不會丟失的。可能出現(xiàn)的問題就是在關(guān)閉的一瞬間,已經(jīng)從消息隊列中取出,但還沒有被處理的消息。
因此我們需要一套平滑關(guān)閉的機制,保證在重啟的時候,已接收的消息可以得到正常處理。
2.問題分析
平滑關(guān)閉的思路如下:
- 在關(guān)閉程序時,首先關(guān)閉消息訂閱,保證不再接收新的消息。
- 關(guān)閉線程池,等待線程池中的消息處理完畢。
- 程序退出。
關(guān)閉消息訂閱:消息隊列的客戶端都會提供關(guān)閉連接的方法,具體可以自行查看API。
關(guān)閉線程池:Java的ThreadPoolExecutor線程池提供shutdown()和shutdownNow()兩個方法,區(qū)別是前者會等待線程池中的消息都處理完畢,后者會直接停止所有線程并返回未處理完的線程List。因為我們需要使用shutdown()方法進行關(guān)閉,并通過isTerminated()方法,判斷線程池是否已經(jīng)關(guān)閉。
那么問題又來了,我們?nèi)绾瓮ㄖ匠绦?,需要?zhí)行關(guān)閉操作呢?
在Linux中,進程的關(guān)閉是通過信號傳遞的,我們可以用kill -9 pid關(guān)閉進程,除了-9之外,我們可以通過 kill -l,查看kill命令的其它信號量。

這里提供兩種關(guān)閉方法:
- 程序中添加
Runtime.getRuntime().addShutdownHook鉤子方法,SIGTERM,SIGINT,SIGHUP三種信號都會觸發(fā)該方法(分別對應(yīng)kill -1/kill -2/kill -15,Ctrl+C也會觸發(fā)SIGINT信號)。 - 程序中通過Signal類注冊信號監(jiān)聽,比如USR2(對應(yīng)kill -12),在handle方法中執(zhí)行關(guān)閉操作。
補充說明:addShutdownHook方法和handle方法中如果再調(diào)用System.exit,會造成deadlock,使進程無法正常退出。
偽代碼分別如下
Runtime.getRuntime().addShutdownHook(new Thread() {
public void run() {
//關(guān)閉訂閱者
//關(guān)閉線程池
//退出
}
});
//注冊linux kill信號量 kill -12
Signal sig = new Signal("USR2");
Signal.handle(sig, new SignalHandler() {
@Override
public void handle(Signal signal) {
//關(guān)閉訂閱者
//關(guān)閉線程池
//退出
}
});
模擬Demo
下面通過一個demo模擬相關(guān)邏輯操作
首先模擬一個生產(chǎn)者,每秒生產(chǎn)5個消息
然后模擬一個訂閱者,收到消息后,放入線程池進行處理,線程池固定4個線程,每個線程處理時間1秒,這樣線程池每秒會積壓1個消息。
package com.lujianing.demo;
import sun.misc.Signal;
import sun.misc.SignalHandler;
import java.util.concurrent.*;
/**
* @author lujianing01@58.com
* @Description:
* @date 2016/11/14
*/
public class MsgClient {
//模擬消費線程池 同時4個線程處理
private static final ThreadPoolExecutor THREAD_POOL = (ThreadPoolExecutor) Executors.newFixedThreadPool(4);
//模擬消息生產(chǎn)任務(wù)
private static final ScheduledExecutorService SCHEDULED_EXECUTOR_SERVICE = Executors.newSingleThreadScheduledExecutor();
//用于判斷是否關(guān)閉訂閱
private static volatile boolean isClose = false;
public static void main(String[] args) throws InterruptedException {
//注冊鉤子方法
Runtime.getRuntime().addShutdownHook(new Thread() {
public void run() {
close();
}
});
BlockingQueue <String> queue = new ArrayBlockingQueue<String>(100);
producer(queue);
consumer(queue);
}
//模擬消息隊列生產(chǎn)者
private static void producer(final BlockingQueue queue){
//每200毫秒向隊列中放入一個消息
SCHEDULED_EXECUTOR_SERVICE.scheduleAtFixedRate(new Runnable() {
public void run() {
queue.offer("");
}
}, 0L, 200L, TimeUnit.MILLISECONDS);
}
//模擬消息隊列消費者 生產(chǎn)者每秒生產(chǎn)5個 消費者4個線程消費1個1秒 每秒積壓1個
private static void consumer(final BlockingQueue queue) throws InterruptedException {
while (!isClose){
getPoolBacklogSize();
//從隊列中拿到消息
final String msg = (String)queue.take();
//放入線程池處理
if(!THREAD_POOL.isShutdown()) {
THREAD_POOL.execute(new Runnable() {
public void run() {
try {
//System.out.println(msg);
TimeUnit.MILLISECONDS.sleep(1000L);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
});
}
}
}
//查看線程池堆積消息個數(shù)
private static long getPoolBacklogSize(){
long backlog = THREAD_POOL.getTaskCount()- THREAD_POOL.getCompletedTaskCount();
System.out.println(String.format("[%s]THREAD_POOL backlog:%s",System.currentTimeMillis(),backlog));
return backlog;
}
private static void close(){
System.out.println("收到kill消息,執(zhí)行關(guān)閉操作");
//關(guān)閉訂閱消費
isClose = true;
//關(guān)閉線程池,等待線程池積壓消息處理
THREAD_POOL.shutdown();
//判斷線程池是否關(guān)閉
while (!THREAD_POOL.isTerminated()) {
try {
//每200毫秒 判斷線程池積壓數(shù)量
getPoolBacklogSize();
TimeUnit.MILLISECONDS.sleep(200L);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
System.out.println("訂閱者關(guān)閉,線程池處理完畢");
}
static {
String osName = System.getProperty("os.name").toLowerCase();
if(osName != null && osName.indexOf("window") == -1) {
//注冊linux kill信號量 kill -12
Signal sig = new Signal("USR2");
Signal.handle(sig, new SignalHandler() {
@Override
public void handle(Signal signal) {
close();
}
});
}
}
}

當我們在服務(wù)上運行時,通過控制臺可以看到相關(guān)的輸出信息,demo中輸出了線程池的積壓消息個數(shù)
java -cp /home/work/lujianing/msg-queue-client/* com.lujianing.demo.MsgClient

另打開一個終端,通過ps命令查看進程號,或者通過nohup啟動Java進程拿到進程id
ps -fe|grep MsgClient

當我們執(zhí)行kill -12 pid的時候 可以看到關(guān)閉業(yè)務(wù)邏輯

3.總結(jié)
其實不單單消息隊列任務(wù),在常見的RPC服務(wù)中也會見到類似的功能,比如58的SCF,在源碼中,也會分別注冊了USR2信號量和addShutdownHook鉤子方法。
在重啟腳本中,首先會發(fā)送kill -12命令,RPC服務(wù)收到信號后會修改Server狀態(tài)為關(guān)閉。接著會發(fā)送kill -15命令,觸發(fā)鉤子方法,關(guān)閉所有的連接。
好了,以上就是這篇文章的全部內(nèi)容了,希望本文的內(nèi)容對大家的學習或者工作具有一定的參考學習價值,如果有疑問大家可以留言交流,謝謝大家對腳本之家的支持。
相關(guān)文章
一文教你利用Stream?API批量Mock數(shù)據(jù)的方法
在日常開發(fā)的過程中我們經(jīng)常會遇到需要mock一些數(shù)據(jù)的場景,比如說?mock?一些接口的返回或者說?mock?一些測試消息用于隊列生產(chǎn)者發(fā)送消息。本文將教你如何通過?Stream?API?批量?Mock?數(shù)據(jù),需要的可以參考一下2022-09-09
FilenameUtils.getName?函數(shù)源碼分析
這篇文章主要為大家介紹了FilenameUtils.getName?函數(shù)源碼分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪2022-09-09
SpringCloud @RefreshScope刷新機制深入探究
RefeshScope這個注解想必大家都用過,在微服務(wù)配置中心的場景下經(jīng)常出現(xiàn),他可以用來刷新Bean中的屬性配置,那大家對他的實現(xiàn)原理了解嗎?它為什么可以做到動態(tài)刷新呢2023-03-03
Java環(huán)境變量的設(shè)置方法(圖文教程)
想要成功配置Java的環(huán)境變量,那肯定就要安裝JDK,才能開始配置的。2013-05-05
教你怎么用java實現(xiàn)客戶端與服務(wù)器一問一答
這篇文章主要介紹了教你怎么用java實現(xiàn)客戶端與服務(wù)器一問一答,文中有非常詳細的代碼示例,對正在學習java的小伙伴們有非常好的幫助,需要的朋友可以參考下2021-04-04
javaWeb用戶權(quán)限控制簡單實現(xiàn)過程
這篇文章主要為大家詳細介紹了javaWeb用戶權(quán)限控制簡單實現(xiàn)過程,具有一定的參考價值,感興趣的小伙伴們可以參考一下2016-08-08
解決SpringBoot運行報錯:找不到或無法加載主類的問題
這篇文章主要介紹了解決SpringBoot運行報錯:找不到或無法加載主類的問題,具有很好的參考價值,對大家的學習或工作有一定的參考價值,需要的朋友可以參考下2023-09-09

