Java中流式并行操作parallelStream的原理和使用方法
Java中流式并行操作parallelStream
0. 問(wèn)題的產(chǎn)生
某天上線后,發(fā)現(xiàn)線上存在一些報(bào)錯(cuò),遂即自己嘗試線上操作,但是發(fā)現(xiàn)功能正常。追蹤相關(guān)的報(bào)錯(cuò)代碼行如下:
PointMissionPO missionPO = missionBO.getDbData();
該行報(bào)錯(cuò)為空指針異常,可以從代碼中判斷唯一能報(bào)出空指針異常的位置為missionBO為空,向上追蹤該引用:
for (PointMissionBO missionBO : result) {
// 循環(huán)體內(nèi)容
}
為一個(gè)列表List的循環(huán)體。于是接著向前追蹤引用,發(fā)現(xiàn)所有處理該列表引用的地方均使用了stream流式操作。經(jīng)常使用函數(shù)式編程的同學(xué)都知道,這很少會(huì)出現(xiàn)null對(duì)象在列表中。我通體檢查了一遍都未發(fā)現(xiàn)任何可能產(chǎn)生null對(duì)象的插入位置。很奇怪那么這個(gè)null對(duì)象是如何被加入到列表中的呢?
后來(lái)有個(gè)小伙伴經(jīng)過(guò)AI上下文分析,給出了可能的位置:
missionByType.entrySet().parallelStream()
.filter(entry -> !CollectionUtils.isEmpty(entry.getValue()))
.map(entry -> processMissions(ctmId, hruId, entry.getValue()))
.forEach(result::addAll);仔細(xì)一看居然使用了parallelStream并行處理,那么什么是parallelStream呢?為什么它會(huì)出現(xiàn)錯(cuò)誤呢?
1. 什么是parallelStream?
Java 8引入的Stream API提供了兩種處理方式:
- stream():串行處理,按順序處理元素
- parallelStream():并行處理,利用多核CPU將數(shù)據(jù)分割成多個(gè)部分并行處理
2. parallelStream的工作原理
parallelStream基于Fork/Join框架實(shí)現(xiàn):
- 將數(shù)據(jù)源分割成多個(gè)子任務(wù)(fork)
- 在不同線程上并行處理這些子任務(wù)
- 合并結(jié)果(join)
3. parallelStream的正確與錯(cuò)誤使用示例
錯(cuò)誤使用示例(來(lái)自我們的測(cè)試代碼):
List<String> result = new ArrayList<>();
missionByType.entrySet().parallelStream()
.map(entry -> processMissions(entry.getValue()))
.forEach(result::addAll); // 危險(xiǎn)操作!
問(wèn)題分析:
- ArrayList不是線程安全的集合
- 多個(gè)線程同時(shí)調(diào)用result::addAll會(huì)產(chǎn)生競(jìng)態(tài)條件
- 可能導(dǎo)致數(shù)據(jù)丟失、重復(fù)、甚至程序崩潰
正確使用方式一:使用同步塊
List<String> result = new ArrayList<>();
missionByType.entrySet().parallelStream()
.map(entry -> processMissions(entry.getValue()))
.filter(processedMissions -> !processedMissions.isEmpty())
.forEach(processedMissions -> {
synchronized (result) {
result.addAll(processedMissions);
}
});
正確使用方式二:使用收集器(推薦)
List<String> safeResult = missionByType.entrySet().parallelStream()
.flatMap(entry -> processMissions(entry.getValue()).stream())
.collect(Collectors.toList());
4. parallelStream在實(shí)際業(yè)務(wù)中的應(yīng)用
查看我們項(xiàng)目中的實(shí)際應(yīng)用案例:
// PointBusinessServiceImpl.java 中的實(shí)際使用
List<PointMissionBO> result = missionByType.entrySet().parallelStream()
.filter(entry -> !CollectionUtils.isEmpty(entry.getValue()))
.flatMap(entry -> processMissions(ctmId, hruId, entry.getValue()).stream())
.collect(Collectors.toList());
這種方式的優(yōu)點(diǎn):
- 使用flatMap展平數(shù)據(jù)結(jié)構(gòu)
- 使用collect收集結(jié)果,避免線程安全問(wèn)題
- 提高了任務(wù)處理效率
5. parallelStream適用場(chǎng)景與注意事項(xiàng)
適用場(chǎng)景:
- 數(shù)據(jù)量較大(通常萬(wàn)級(jí)以上)
- 計(jì)算密集型操作
- 無(wú)狀態(tài)操作(函數(shù)式編程)
- 不依賴處理順序的操作
不適用場(chǎng)景:
- 數(shù)據(jù)量?。ú⑿虚_(kāi)銷可能超過(guò)收益)
- IO密集型操作
- 有狀態(tài)共享操作
- 需要保證處理順序的場(chǎng)景
注意事項(xiàng):
- 線程安全:避免在并行流中使用非線程安全的對(duì)象
- 副作用:避免在流操作中修改外部狀態(tài)
- 性能考量:并行不一定比串行快,需根據(jù)實(shí)際情況評(píng)估
- 資源競(jìng)爭(zhēng):注意共享資源的訪問(wèn)控制
6. 最佳實(shí)踐建議
- 優(yōu)先考慮collect:使用收集器而不是直接修改共享集合
- 避免副作用:確保流操作是無(wú)狀態(tài)的純函數(shù)
- 合理選擇數(shù)據(jù)結(jié)構(gòu):使用適合并行處理的數(shù)據(jù)結(jié)構(gòu)
- 測(cè)試性能:在實(shí)際環(huán)境中測(cè)試并行處理效果
- 監(jiān)控資源使用:關(guān)注CPU和內(nèi)存使用情況
以上就是關(guān)于parallelStream并行處理的AI輸出內(nèi)容。下面結(jié)合實(shí)際問(wèn)題發(fā)生的場(chǎng)景進(jìn)行推演和論述。
7. 其他思考
上面說(shuō)了這么多,核心重點(diǎn)在于:parallelStream中使用線程不安全的對(duì)象操作時(shí)會(huì)出現(xiàn)異常。那么具體是什么樣的異常呢?
其實(shí)我們仔細(xì)想想非內(nèi)存安全List的內(nèi)存管理,熟悉八股文的同學(xué)應(yīng)該一下子就懂了,沒(méi)錯(cuò)就是內(nèi)存擴(kuò)展機(jī)制。
當(dāng)一個(gè)非內(nèi)存安全的List在觸發(fā)到內(nèi)存擴(kuò)展閾值的時(shí)候,就會(huì)觸發(fā)一次內(nèi)存擴(kuò)展。具體原理就是開(kāi)辟一個(gè)更長(zhǎng)的(通常是2倍)列表,然后將原數(shù)據(jù)賦值到新列表中。這個(gè)過(guò)程中,如果出現(xiàn)了并行開(kāi)辟的情況,那賦值的內(nèi)容以及目標(biāo)列表就會(huì)變得混亂,出現(xiàn)null對(duì)象也并非不可能。因此當(dāng)嘗試指定初始列表大小為256的情況下,第0章的錯(cuò)誤就自然消失了。
那么也就出現(xiàn)一些關(guān)于Collection的使用建議:
- 初始化List時(shí),最好能夠有一個(gè)預(yù)估的列表大小并指定,其實(shí)所有Collection對(duì)象都可以這么做。
- 但凡出現(xiàn)非線程安全的Collection對(duì)象需要參與并行計(jì)算時(shí),都需要注意它的數(shù)據(jù)正確性應(yīng)當(dāng)如何保證,如果無(wú)法自行控制,或者控制有難度的,可以考慮使用Concurrent包中的內(nèi)容。
到此這篇關(guān)于Java中流式并行操作parallelStream的原理和使用方法的文章就介紹到這了,更多相關(guān)java并行parallelStream內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Spring5中SpringWebContext方法過(guò)時(shí)的解決方案
這篇文章主要介紹了Spring5中SpringWebContext方法過(guò)時(shí)的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2022-01-01
java合成模式之神奇的樹(shù)結(jié)構(gòu)
這篇文章主要介紹了java合成模式,文中運(yùn)用大量的代碼進(jìn)行詳細(xì)講解,希望大家看完本文后能學(xué)習(xí)到相關(guān)的知識(shí),需要的朋友可以參考一下2021-08-08
MyBatis框架關(guān)聯(lián)映射實(shí)例詳解
這篇文章主要介紹了MyBatis框架關(guān)聯(lián)映射,關(guān)系映射主要處理復(fù)雜的SQl查詢,如子查詢,多表聯(lián)查等復(fù)雜查詢,應(yīng)用此種需求時(shí)可以考慮使用,需要的朋友可以參考下2022-11-11
簡(jiǎn)單講解Java設(shè)計(jì)模式編程中的單一職責(zé)原則
這篇文章主要介紹了Java設(shè)計(jì)模式編程中的單一職責(zé)原則,這在團(tuán)隊(duì)開(kāi)發(fā)編寫(xiě)接口時(shí)經(jīng)常使用這樣的約定,需要的朋友可以參考下2016-02-02
SpringBoot整合log4j日志與HashMap的底層原理解析
這篇文章主要介紹了SpringBoot整合log4j日志與HashMap的底層原理,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-01-01
java實(shí)現(xiàn)簡(jiǎn)單的搜索引擎
這篇文章主要為大家詳細(xì)介紹了java實(shí)現(xiàn)簡(jiǎn)單的搜索引擎的相關(guān)資料,需要的朋友可以參考下2016-02-02
java使用CompletableFuture分批處理任務(wù)實(shí)現(xiàn)
本文主要介紹了java使用CompletableFuture分批處理任務(wù)實(shí)現(xiàn),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧2024-07-07
MybatisPlus修改時(shí)空字段無(wú)法修改的解決方案
這篇文章主要介紹了MybatisPlus修改時(shí)空字段無(wú)法修改的解決方案,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-09-09

