공공부하자개발 · 영어 학습 노트
자바
고급모던 자바와 성능0/10 완료
  • 01제네릭과 와일드카드
  • 02멀티스레드와 동기화
  • 03람다식과 함수형 인터페이스
  • 04Stream API와 병렬 처리
  • 05Optional로 NPE 방지
  • 06메서드 활용 패턴 (고급)
  • 07Java 21 모던 문법
  • 08어노테이션·리플렉션·동적 프록시
  • 09CompletableFuture 심화와 가상 스레드 실전
  • 10JVM 메모리·GC·OOM 진단
사이트 소개개인정보처리방침연락처
© 2026 공부하자
홈 › 고급 › 04 / 10

Stream API와 병렬 처리

섹션 7진행 0 / 10
1왜 배우는가2핵심 원리3코드 예제4응용 변형 예제5자주 하는 실수 (Tip)6연습 문제7정리‹ 이전다음 ›

3. 코드 예제

예제 1: 지연 평가와 단락 평가 증명

2.3 의 설명을 실행 가능한 코드로 확인합니다. sorted 가 장벽인 것도 함께 봅니다.

java
import java.util.*;
import java.util.stream.*;

public class LazyProof {
    public static void main(String[] args) {
        Stream<String> pipeline = Stream.of("a", "bb", "ccc", "dddd")
            .peek(x -> System.out.println("filter 전: " + x))
            .filter(x -> x.length() >= 2)
            .peek(x -> System.out.println("  map 전: " + x))
            .map(String::toUpperCase);
        System.out.println("--- 중간 연산만 쌓음, 아직 실행 안 됨 ---");
        System.out.println(pipeline.limit(2).toList());

        System.out.println("--- sorted 는 전부 모은 뒤에야 다음으로 넘긴다 ---");
        Stream.of(3, 1, 2)
            .peek(x -> System.out.println("sorted 전: " + x))
            .sorted()
            .peek(x -> System.out.println("  sorted 후: " + x))
            .findFirst();

        System.out.println("--- 무한 스트림 + takeWhile ---");
        System.out.println(Stream.iterate(1, x -> x * 2).takeWhile(x -> x < 100).toList());
    }
}
// 출력:
// --- 중간 연산만 쌓음, 아직 실행 안 됨 ---
// filter 전: a
// filter 전: bb
//   map 전: bb
// filter 전: ccc
//   map 전: ccc
// [BB, CCC]
// --- sorted 는 전부 모은 뒤에야 다음으로 넘긴다 ---
// sorted 전: 3
// sorted 전: 1
// sorted 전: 2
//   sorted 후: 1
// --- 무한 스트림 + takeWhile ---
// [1, 2, 4, 8, 16, 32, 64]

예제 2: 주문 집계 — 고객별 매출 상위 N, 월별 건수, 카테고리별 평균

실무에서 가장 흔한 세 가지 집계입니다. groupingBy + downstream 컬렉터 조합을 봅니다.

java
import java.time.*;
import java.util.*;
import java.util.stream.*;
import static java.util.stream.Collectors.*;

public class OrderAggregation {
    record Order(String id, String customer, String category, long amount, LocalDate date) {}

    static List<Order> sample() {
        return List.of(
            new Order("O1", "kim", "book", 12_000, LocalDate.of(2024, 1, 5)),
            new Order("O2", "lee", "food", 8_000,  LocalDate.of(2024, 1, 20)),
            new Order("O3", "kim", "food", 15_000, LocalDate.of(2024, 2, 3)),
            new Order("O4", "park", "book", 30_000, LocalDate.of(2024, 2, 14)),
            new Order("O5", "lee", "book", 5_000,  LocalDate.of(2024, 3, 1)),
            new Order("O6", "kim", "toy", 22_000,  LocalDate.of(2024, 3, 9))
        );
    }

    public static void main(String[] args) {
        List<Order> orders = sample();

        // 1) 고객별 매출 합계 → 상위 2명
        Map<String, Long> byCustomer = orders.stream()
            .collect(groupingBy(Order::customer, summingLong(Order::amount)));
        List<Map.Entry<String, Long>> top2 = byCustomer.entrySet().stream()
            .sorted(Map.Entry.<String, Long>comparingByValue().reversed())
            .limit(2).toList();
        System.out.println("top2 = " + top2);

        // 2) 월별 주문 건수 (TreeMap 으로 월 순서 정렬)
        Map<Month, Long> byMonth = orders.stream()
            .collect(groupingBy(o -> o.date().getMonth(), TreeMap::new, counting()));
        System.out.println("byMonth = " + byMonth);

        // 3) 카테고리별 평균 금액
        Map<String, Double> avgByCategory = orders.stream()
            .collect(groupingBy(Order::category, TreeMap::new, averagingLong(Order::amount)));
        System.out.println("avgByCategory = " + avgByCategory);

        // 4) 월별 → 카테고리별 매출 (2단 그룹핑)
        Map<Month, Map<String, Long>> nested = orders.stream()
            .collect(groupingBy(o -> o.date().getMonth(), TreeMap::new,
                     groupingBy(Order::category, TreeMap::new, summingLong(Order::amount))));
        System.out.println("nested = " + nested);

        // 5) 통계 한 번에
        LongSummaryStatistics stats = orders.stream().mapToLong(Order::amount).summaryStatistics();
        System.out.printf("count=%d sum=%d avg=%.1f max=%d%n", stats.getCount(), stats.getSum(), stats.getAverage(), stats.getMax());
    }
}
// 출력:
// top2 = [kim=49000, lee=13000]
// byMonth = {JANUARY=2, FEBRUARY=2, MARCH=2}
// avgByCategory = {book=15666.666666666666, food=11500.0, toy=22000.0}
// nested = {JANUARY={book=12000, food=8000}, FEBRUARY={book=30000, food=15000}, MARCH={book=5000, toy=22000}}
// count=6 sum=92000 avg=15333.3 max=30000

예제 3: toMap 키 충돌, partitioningBy, joining, teeing, mapping

컬렉터의 세부 동작을 한 번에 확인합니다. 특히 toMap 이 키 충돌에서 어떻게 터지고 어떻게 막는지 봅니다.

java
import java.util.*;
import java.util.stream.*;
import static java.util.stream.Collectors.*;

public class CollectorsDemo {
    record Order(String id, String customer, long amount, boolean paid) {}

    public static void main(String[] args) {
        List<Order> orders = List.of(
            new Order("O1", "kim", 100, true), new Order("O2", "lee", 200, false),
            new Order("O3", "kim", 300, true), new Order("O4", "park", 50, true));

        // toMap 키 충돌
        try {
            orders.stream().collect(toMap(Order::customer, Order::amount));
        } catch (IllegalStateException e) {
            System.out.println("충돌: " + e.getMessage());
        }
        Map<String, Long> merged = orders.stream()
            .collect(toMap(Order::customer, Order::amount, Long::sum, TreeMap::new));   // 합산 + 정렬
        System.out.println("merged = " + merged);

        // partitioningBy: 결제/미결제 (항상 true/false 두 키)
        Map<Boolean, List<String>> paidSplit = orders.stream()
            .collect(partitioningBy(Order::paid, mapping(Order::id, toList())));
        System.out.println("paid = " + paidSplit.get(true) + ", unpaid = " + paidSplit.get(false));

        // joining
        String csv = orders.stream().map(Order::id).collect(joining(",", "[", "]"));
        System.out.println("csv = " + csv);

        // teeing: 최소·최대를 한 번의 순회로
        String range = orders.stream().collect(teeing(
            minBy(Comparator.comparingLong(Order::amount)),
            maxBy(Comparator.comparingLong(Order::amount)),
            (min, max) -> min.get().amount() + "~" + max.get().amount()));
        System.out.println("range = " + range);

        // groupingBy + mapping + collectingAndThen: 고객별 주문 id 를 불변 Set 으로
        Map<String, Set<String>> idsByCustomer = orders.stream()
            .collect(groupingBy(Order::customer, TreeMap::new,
                     mapping(Order::id, collectingAndThen(toSet(), Collections::unmodifiableSet))));
        System.out.println("ids = " + idsByCustomer);
    }
}
// 출력:
// 충돌: Duplicate key kim (attempted merging values 100 and 300)
// merged = {kim=400, lee=200, park=50}
// paid = [O1, O3, O4], unpaid = [O2]
// csv = [O1,O2,O3,O4]
// range = 50~300
// ids = {kim=[O1, O3], lee=[O2], park=[O4]}

예제 4: flatMap 으로 중첩 리스트 펼치기 + 대용량 로그 파싱

주문 → 상품 라인을 펼쳐 상품별 판매 수량을 집계하고, 로그 파일을 Files.lines 로 읽어 레벨별 건수와 ERROR 메시지 상위를 뽑습니다.

java
import java.io.IOException;
import java.nio.file.*;
import java.util.*;
import java.util.stream.*;
import static java.util.stream.Collectors.*;

public class FlatMapAndLogs {
    record Line(String sku, int qty) {}
    record Order(String id, List<Line> lines) {}

    public static void main(String[] args) throws IOException {
        List<Order> orders = List.of(
            new Order("O1", List.of(new Line("pen", 2), new Line("book", 1))),
            new Order("O2", List.of(new Line("pen", 5))),
            new Order("O3", List.of(new Line("book", 3), new Line("cup", 1))));

        // map 이면 Stream<List<Line>>, flatMap 이면 Stream<Line>
        Map<String, Integer> qtyBySku = orders.stream()
            .flatMap(o -> o.lines().stream())
            .collect(groupingBy(Line::sku, TreeMap::new, summingInt(Line::qty)));
        System.out.println("qtyBySku = " + qtyBySku);

        // 로그 파일 생성 (실무에선 수 GB 파일도 lines() 는 한 줄씩 읽으므로 메모리 안전)
        Path log = Files.createTempFile("app", ".log");
        Files.write(log, List.of(
            "2024-03-01 10:00:01 INFO  server started",
            "2024-03-01 10:00:05 ERROR db connection refused",
            "2024-03-01 10:00:09 WARN  slow query 1200ms",
            "2024-03-01 10:01:00 ERROR db connection refused",
            "2024-03-01 10:02:00 ERROR payment timeout",
            "malformed line"));

        try (Stream<String> lines = Files.lines(log)) {           // 반드시 닫기
            Map<String, Long> byLevel = lines
                .map(l -> l.split("\\s+", 4))
                .filter(p -> p.length == 4)                        // 깨진 줄 제외
                .collect(groupingBy(p -> p[2], TreeMap::new, counting()));
            System.out.println("byLevel = " + byLevel);
        }
        try (Stream<String> lines = Files.lines(log)) {
            List<Map.Entry<String, Long>> topErrors = lines
                .filter(l -> l.contains(" ERROR "))
                .map(l -> l.substring(l.indexOf("ERROR") + 6))
                .collect(groupingBy(m -> m, counting()))
                .entrySet().stream()
                .sorted(Map.Entry.<String, Long>comparingByValue().reversed().thenComparing(Map.Entry.comparingByKey()))
                .limit(2).toList();
            System.out.println("topErrors = " + topErrors);
        }
        Files.delete(log);
    }
}
// 출력:
// qtyBySku = {book=4, cup=1, pen=7}
// byLevel = {ERROR=3, INFO=1, WARN=1}
// topErrors = [db connection refused=2, payment timeout=1]

예제 5: 병렬 스트림 — 이득/손해, 공유 상태 오염, forEachOrdered

CPU 바운드 작업으로 병렬 이득을 확인하고, 공유 ArrayList 를 쓰면 결과가 깨지는 것과 collect 로 고치는 것을 보여줍니다.

java
import java.util.*;
import java.util.concurrent.*;
import java.util.stream.*;

public class ParallelDemo {
    static boolean isPrime(int n) {
        if (n < 2) return false;
        for (int i = 2; (long) i * i <= n; i++) if (n % i == 0) return false;
        return true;
    }

    public static void main(String[] args) throws Exception {
        int limit = 2_000_000;
        long t0 = System.nanoTime();
        long seq = IntStream.range(0, limit).filter(ParallelDemo::isPrime).count();
        long t1 = System.nanoTime();
        long par = IntStream.range(0, limit).parallel().filter(ParallelDemo::isPrime).count();
        long t2 = System.nanoTime();
        System.out.println("primes = " + seq + " / " + par);
        System.out.println("parallel faster: " + ((t2 - t1) < (t1 - t0)));   // 코어 2개 이상이면 true

        // 공유 상태 오염: ArrayList 는 스레드 안전하지 않음
        List<Integer> shared = new ArrayList<>();
        try {
            IntStream.range(0, 100_000).parallel().forEach(shared::add);
            System.out.println("shared size = " + shared.size() + (shared.size() == 100_000 ? "" : "  <-- 손실"));
        } catch (ArrayIndexOutOfBoundsException e) {
            System.out.println("shared list 내부 배열 깨짐: " + e.getClass().getSimpleName());
        }
        List<Integer> safe = IntStream.range(0, 100_000).parallel().boxed().toList();   // collect 는 안전
        System.out.println("safe size = " + safe.size() + ", ordered = " + (safe.get(99_999) == 99_999));

        // forEach 는 순서 무작위, forEachOrdered 는 인카운터 순서
        StringBuilder ordered = new StringBuilder();
        IntStream.rangeClosed(1, 5).parallel().forEachOrdered(ordered::append);
        System.out.println("forEachOrdered = " + ordered);

        // 커스텀 풀에서 실행
        ForkJoinPool pool = new ForkJoinPool(2);
        long sum = pool.submit(() -> LongStream.rangeClosed(1, 1_000_000).parallel().sum()).get();
        pool.shutdown();
        System.out.println("sum in custom pool = " + sum);
    }
}
// 출력 (shared 줄은 실행마다 다름):
// primes = 148933 / 148933
// parallel faster: true
// shared size = 97412  <-- 손실
// safe size = 100000, ordered = true
// forEachOrdered = 12345
// sum in custom pool = 500000500000

예제 직접 실행

아래 폴더를 JDK 21 로 컴파일하고 실행합니다.

cd java-src\advanced\04_stream
javac -encoding UTF-8 *.java && java Main
코드 예제
  • 예제 1: 지연 평가와 단락 평가 증명
  • 예제 2: 주문 집계 — 고객별 매출 상위 N, 월별 건수, 카테고리별 평균
  • 예제 3: toMap 키 충돌, partitioningBy, joining, teeing, mapping
  • 예제 4: flatMap 으로 중첩 리스트 펼치기 + 대용량 로그 파싱
  • 예제 5: 병렬 스트림 — 이득/손해, 공유 상태 오염, forEachOrdered
이전 섹션2 핵심 원리3 / 7다음 섹션4 응용 변형 예제