为什么我的Java Stream流操作会吃掉内存?
为ä»ä¹æçJava Streamæµæä½ä¼åæå åï¼
é¿æ©çç¾å®ç®± 2026-08-09 0 é 读1åé- 为ä»ä¹æçJava Streamæµæä½ä¼åæå åï¼*
å¼è¨
Java 8å¼å ¥çStream APIæå¤§å°ç®åäºéåæä½ï¼æä¾äºå£°æå¼ã彿°å¼çç¼ç¨é£æ ¼ãç¶èï¼è®¸å¤å¼åè å¨å®é 使ç¨ä¸åç°ï¼æäºStreamæä½ä¼æå¤æ¶è大éå åï¼çè³å¯¼è´OutOfMemoryErrorãè¿ç§ç°è±¡èåçåå æ¯ä»ä¹ï¼å¦ä½é¿å Streamæä½æä¸ºå åææï¼æ¬æå°æ·±å ¥åæStreamçå åä½¿ç¨æºå¶ï¼æç¤ºå¸¸è§é·é±ï¼å¹¶æä¾ä¼å建议ã
ä¸ãStreamçå 忍¡ååºç¡
1.1 Streamçæ°æ§æ±å¼ä¸ä¸é´æä½
Streamæä½å为ä¸é´æä½ï¼Intermediate Operationsï¼åç»æ¢æä½ï¼Terminal Operationsï¼ãå ³é®ç¹æ§å¨äºæ°æ§æ±å¼ââåªæéå°ç»æ¢æä½æ¶æä¼è§¦åå®é 计ç®ãè¿ç§è®¾è®¡ç论ä¸å¯ä»¥èçå åï¼ä½å¨æäºåºæ¯ä¸ä¼äº§çç¸åææã
List<Integer> list = Arrays.asList(1, 2, 3, 4);
Stream<Integer> stream = list.stream()
.filter(x -> x > 1) // ä¸é´æä½ï¼ä¸ç«å³æ§è¡
.map(x -> x * 2); // ä¸é´æä½
1.2 æµæ°´çº¿æ§è¡æºå¶
å½ç»æ¢æä½è¢«è°ç¨æ¶ï¼JVMä¼æå»ºä¸ä¸ªæä½æµæ°´çº¿ï¼Operation Pipelineï¼ãè¿ä¸ªæµæ°´çº¿ä¸æ¯ç®åç彿°ç»åï¼èæ¯éè¿Sinkæ¥å£å®ç°ç夿ååéä¿¡æºå¶ãæ¯ä¸ªä¸é´æä½é½ä¼å¨å åä¸å建对åºçStage对象ã
äºãå åæ¶èçåå¤§æ ¹æº
2.1 ä¸é´ç¶æçç©åï¼Materializationï¼
æäºæä½ä¼å¼ºå¶ç©åä¸é´ç»æï¼
- sorted()ï¼å¿ 须尿æå ç´ æ¶éå°å åæè½æåº
// æ¶èO(n)å
å
List<Integer> sorted = largeList.stream()
.sorted()
.collect(Collectors.toList());
- distinct()ï¼éè¦ç»´æ¤åå¸è¡¨è®°å½å·²æå ç´
// å
åæ¶èåå³äºå¯ä¸å
ç´ æ°é
List<Integer> unique = largeList.stream()
.distinct()
.collect(Collectors.toList());
2.2 è£ ç®±/æç®±å¼é
åå§ç±»åæµï¼IntStreamçï¼ä¸å¯¹è±¡æµï¼Streamï¼ä¹é´ç转æ¢ï¼
// æ¯ä¸ªint被è£
箱为Integer
IntStream.range(0, 1_000_000)
.boxed() // 产ç大éInteger对象
.collect(Collectors.toList());
2.3 å¹¶è¡æµç线ç¨å¼é
å¹¶è¡æµï¼parallelStreamï¼ä½¿ç¨ForkJoinPoolï¼
- 任塿å产çé¢å¤å¯¹è±¡
- å·¥ä½éåå ç¨å å
// æ¯ä¸ªçº¿ç¨éè¦ç»´æ¤èªå·±çä¸é´ç»æ
List<Integer> result = largeList.parallelStream()
.filter(x -> x % 2 == 0)
.collect(Collectors.toList());
2.4 ç»ç»æä½çæ¶éå¨é®é¢
æäºCollectorsä¼ä¿çä¸é´ç¶æï¼
- groupingByï¼ç»´æ¤å®æ´çMapç»æ
// 妿åç»é®å¾å¤ï¼Mapä¼åå¾å·¨å¤§
Map<Integer, List<Item>> groups = items.stream()
.collect(Collectors.groupingBy(Item::getCategory));
- joiningï¼StringBuilderæç»å¢é¿
// 大éåæ¼æ¥å¯è½å¯¼è´è¶
大StringBuilder
String combined = largeList.stream()
.map(Object::toString)
.collect(Collectors.joining(","));
ä¸ãæ§è½é·é±ä¸ä¼åæ¹æ¡
3.1 æåºä¼åçç¥
- é®é¢åºæ¯*ï¼
// å®å
¨ç©å两个å表
List<Item> result = items.stream()
.sorted(comparing(Item::getPrice))
.sorted(comparing(Item::getWeight))
.collect(Collectors.toList());
- è§£å³æ¹æ¡*ï¼
- 使ç¨å¤åæ¯è¾å¨
List<Item> result = items.stream()
.sorted(comparing(Item::getPrice)
.thenComparing(Item::getWeight))
.collect(Collectors.toList());
- 对äºå¤§æ°æ®éèèä½¿ç¨æ°æ®åºæåº
3.2 åå§ç±»åæµä¼å
- 使忳*ï¼
Stream<Integer> stream = IntStream.range(0, 1_000_000)
.boxed();
- 髿æ¿ä»£*ï¼
IntStream stream = IntStream.range(0, 1_000_000);
// 使ç¨ä¸ç¨æ¹æ³å¦sum(), average()
3.3 å¹¶è¡æµä½¿ç¨åå
- ä¸éç¨åºæ¯*ï¼
- æ°æ®éå°ï¼<10,000å ç´ ï¼
- æIOé»å¡æä½
- ä¾èµé¡ºåºçæä½ï¼å¦limit, findFirstï¼
- æä½³å®è·µ*ï¼
// æç¡®è®¾ç½®å¹¶è¡éå¼
List<Integer> result = largeList.parallelStream()
.filter(x -> expensiveCalculation(x))
.collect(Collectors.toList());
3.4 å ååå¥½çæ¶éå¨
- æ¿ä»£groupingBy*ï¼
// 使ç¨ä¸æ¸¸æ¶é卿§å¶å
å
Map<Integer, Long> countByCategory = items.stream()
.collect(groupingBy(
Item::getCategory,
Collectors.counting() // ä¸ä¿åå
¨é¨å
ç´
));
- å¤§ææ¬æ¼æ¥ä¼å*ï¼
// 使ç¨StringWriterç¼å²å°æä»¶
StringWriter writer = new StringWriter();
items.stream()
.map(Object::toString)
.forEach(writer::write);
åãè¯æå·¥å ·ä¸ææ¯
4.1 å ååæå·¥å ·
- VisualVM/JConsoleï¼è§å¯å å ååå
- JFRï¼Java Flight Recorderï¼ï¼æè·Streamç¸å ³äºä»¶
- Heap Dumpåæï¼è¯å«ä¿ççä¸é´éå
4.2 åºåæµè¯æ¹æ³
使ç¨JMHè¿è¡ç²¾ç¡®æµéï¼
@Benchmark
public void testStreamMemory(Blackhole bh) {
List<Integer> result = IntStream.range(0, 100_000)
.boxed()
.filter(x -> x % 2 == 0)
.collect(Collectors.toList());
bh.consume(result);
}
4.3 JVMåæ°è°ä¼
ç¸å ³åæ°ï¼
- XX:+PrintGCDetails # è§å¯GCè¡ä¸º
- XX:NativeMemoryTracking=detail # è·è¸ªå
é¨å
å
äºãæ¶æçº§è§£å³æ¹æ¡
5.1 åæ¹å¤ç模å¼
// 使ç¨Streamçsplititerator
Spliterator<Item> spliterator = items.spliterator();
StreamSupport.stream(spliterator, false)
.batch(1000) // èªå®ä¹åæ¹
.forEach(this::processBatch);
5.2 å¤é¨æåºæ¿ä»£æ¹æ¡
对äºè¶ å¤§æ°æ®éï¼
// 使ç¨åEhcacheè¿æ ·çç¼åç³»ç»
Cache<Integer, Item> cache = ...;
items.stream()
.sorted(externalComparator)
.forEach(item -> cache.put(item.getId(), item));
5.3 ååºå¼ç¼ç¨æ´å
ç»åProject Reactorå¤çèåï¼
Flux.fromIterable(items)
.filter(x -> x > 0)
.subscribeOn(Schedulers.parallel())
.subscribe(this::process);
æ»ç»
Java Streamçå åæ¶èé®é¢å¾å¾æºäºå¯¹åºå±æºå¶çä¸å®å ¨çè§£ãéè¿è¯å«ç©åæä½ãä¼åæ¶éå¨ä½¿ç¨ãåçéæ©åå§ç±»åæµä»¥åè°¨æ 使ç¨å¹¶è¡æµï¼å¯ä»¥æ¾èéä½å åå ç¨ãè®°ä½ï¼Stream APIçç®æ´æ§èåéèçå®ç°å¤ææ§ï¼æ§è½ä¼åéè¦åºäºå®é æµéèéå设ãå¨å¤§æ°æ®åºæ¯ä¸ï¼èèç»ååæ¹å¤çæååºå¼ç¼ç¨æ¨¡åæè½çæ£åæ¥Streamçä¼å¿ã
Aitishiku.com