공공부하자개발 · 영어 학습 노트
자바
고급모던 자바와 성능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정리‹ 이전다음 ›

4. 응용 변형 예제

3절이 "스트림으로 무엇을 만들 수 있는가"였다면, 여기서는 같은 문제를 for 문·순차·병렬로 나란히 재고, 컬렉터를 직접 정의해 재사용 유틸로 뽑고, 대용량 로그와 null·빈 스트림 같은 엣지 케이스에 적용해 봅니다. 시간이 찍히는 예제는 실행 환경에 따라 값이 달라집니다.

변형 1: 같은 집계를 for 문 → 스트림 → 병렬 스트림으로

예제 2 의 groupingBy + summingLong 을 200만 건에 대해 세 방식으로 구현하고 시간을 잽니다. for 문과 순차 스트림은 같은 일을 하며 차이는 가독성입니다. 병렬은 groupingByConcurrent 로 합치기 비용을 줄였을 때만 이득이고, 데이터가 작거나 CPU 가 이미 바쁘면 오히려 느려질 수 있으므로 반드시 측정 후 택합니다.

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

public class AggregateThreeWays {
    record Order(int customer, long amount) {}

    static Map<Integer, Long> withFor(List<Order> orders) {
        Map<Integer, Long> m = new HashMap<>();
        for (Order o : orders) m.merge(o.customer(), o.amount(), Long::sum);
        return m;
    }
    static Map<Integer, Long> withStream(List<Order> orders) {
        return orders.stream().collect(groupingBy(Order::customer, summingLong(Order::amount)));
    }
    static Map<Integer, Long> withParallel(List<Order> orders) {
        return orders.parallelStream().collect(groupingByConcurrent(Order::customer, summingLong(Order::amount)));
    }

    static long time(String label, java.util.function.Supplier<Map<Integer, Long>> job) {
        long t0 = System.nanoTime();
        Map<Integer, Long> r = job.get();
        long ms = (System.nanoTime() - t0) / 1_000_000;
        System.out.printf("%-9s customers=%d total=%d  %dms%n", label, r.size(), r.values().stream().mapToLong(Long::longValue).sum(), ms);
        return ms;
    }

    public static void main(String[] args) {
        Random rnd = new Random(42);
        List<Order> orders = IntStream.range(0, 2_000_000)
            .mapToObj(i -> new Order(rnd.nextInt(1000), rnd.nextInt(100_000))).toList();

        for (int i = 0; i < 2; i++) { withFor(orders); withStream(orders); withParallel(orders); }   // 워밍업(JIT)
        time("for", () -> withFor(orders));
        time("stream", () -> withStream(orders));
        time("parallel", () -> withParallel(orders));
        System.out.println("same result: " + withFor(orders).equals(withParallel(orders)));
    }
}
// 출력 (시간 값은 환경에 따라 다름, 결과 값은 시드 고정이라 동일):
// for       customers=1000 total=100031088655  729ms
// stream    customers=1000 total=100031088655  831ms
// parallel  customers=1000 total=100031088655  437ms
// same result: true

변형 2: Collector.of 로 만드는 상위 N 컬렉터와 제네릭 topNBy 유틸

예제 2 의 "상위 2명"은 전체를 sorted 한 뒤 limit 했습니다. 100만 건에서 상위 3개만 필요하다면 전부 정렬하는 것은 낭비입니다.

Collector.of(supplier, accumulator, combiner, finisher) 로 크기 N 의 힙만 유지하는 컬렉터를 직접 정의하면 메모리는 N 에 비례하고, combiner 덕분에 병렬에서도 그대로 동작합니다. 이를 groupingBy 의 downstream 으로 꽂아 "분류별 상위 N" 제네릭 유틸로 재사용합니다.

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

public class TopNCollector {
    record Order(String customer, String category, long amount) {}

    // 전체를 정렬하지 않고 크기 n 의 힙(PriorityQueue)만 유지하는 상위 N 컬렉터. 병렬에서도 combiner 로 합쳐진다
    static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> cmp) {
        return Collector.of(
            () -> new PriorityQueue<T>(cmp),                         // supplier: 가장 작은 것이 head
            (pq, t) -> { pq.offer(t); if (pq.size() > n) pq.poll(); }, // accumulator: n 개 초과 시 최소 제거
            (a, b) -> { b.forEach(t -> { a.offer(t); if (a.size() > n) a.poll(); }); return a; }, // combiner
            pq -> { List<T> out = new ArrayList<>(pq); out.sort(cmp.reversed()); return out; }    // finisher
        );
    }

    // 분류 + 상위 N 을 한 줄로 조합하는 제네릭 유틸
    static <T, K> Map<K, List<T>> topNBy(Collection<T> items, java.util.function.Function<T, K> key,
                                         int n, Comparator<? super T> cmp) {
        return items.stream().collect(groupingBy(key, TreeMap::new, topN(n, cmp)));
    }

    public static void main(String[] args) {
        List<Order> orders = List.of(
            new Order("kim", "book", 12_000), new Order("lee", "food", 8_000), new Order("kim", "food", 15_000),
            new Order("park", "book", 30_000), new Order("lee", "book", 5_000), new Order("kim", "toy", 22_000),
            new Order("choi", "food", 9_000));
        Comparator<Order> byAmount = Comparator.comparingLong(Order::amount);

        System.out.println(orders.stream().collect(topN(3, byAmount)).stream().map(Order::amount).toList());
        System.out.println(orders.parallelStream().collect(topN(3, byAmount)).stream().map(Order::amount).toList());

        Map<String, List<Order>> top2ByCategory = topNBy(orders, Order::category, 2, byAmount);
        top2ByCategory.forEach((k, v) -> System.out.println(k + " -> " + v.stream().map(o -> o.customer() + ":" + o.amount()).toList()));

        System.out.println(Stream.<Order>empty().collect(topN(3, byAmount)));   // 빈 스트림도 빈 리스트
    }
}
// 출력:
// [30000, 22000, 15000]
// [30000, 22000, 15000]
// book -> [park:30000, kim:12000]
// food -> [kim:15000, choi:9000]
// toy -> [kim:22000]
// []

변형 3: 30만 줄 로그 파일 스트리밍 파싱

예제 4 는 6줄짜리 로그였습니다. 여기서는 30만 줄(약 12MB)을 생성해 Files.lines 로 두 번 훑으며 레벨별 건수와 상위 에러를 뽑습니다. 파싱 결과를 Optional<Entry> 로 돌려주고 flatMap(Optional::stream) 으로 깨진 줄을 버리는 패턴, 그리고 종단 연산이 보관하는 것이 "레벨 수 · 에러 메시지 종류 수" 뿐이라 파일 크기와 무관하게 힙이 일정하다는 점이 핵심입니다.

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

public class BigLogParse {
    record Entry(String level, String message) {}

    static Optional<Entry> parse(String line) {           // 깨진 줄은 empty
        String[] p = line.split(" ", 4);
        return p.length == 4 ? Optional.of(new Entry(p[2], p[3])) : Optional.empty();
    }

    public static void main(String[] args) throws IOException {
        Path log = Files.createTempFile("app", ".log");
        String[] levels = {"INFO", "INFO", "INFO", "WARN", "ERROR"};
        String[] errors = {"db connection refused", "payment timeout", "cache miss storm"};
        Random rnd = new Random(7);
        try (BufferedWriter w = Files.newBufferedWriter(log)) {           // 30만 줄 생성 (~12MB)
            for (int i = 0; i < 300_000; i++) {
                String lv = levels[rnd.nextInt(levels.length)];
                String msg = lv.equals("ERROR") ? errors[rnd.nextInt(errors.length)] : "request " + i + " ok";
                w.write("2024-03-01 10:00:" + (i % 60) + " " + lv + " " + msg);
                w.newLine();
                if (i % 50_000 == 0) { w.write("malformed"); w.newLine(); }
            }
        }
        System.out.println("file size = " + Files.size(log) / 1024 / 1024 + " MB");

        long t0 = System.nanoTime();
        Map<String, Long> byLevel;
        List<Map.Entry<String, Long>> topErrors;
        try (Stream<String> lines = Files.lines(log)) {                     // 한 줄씩 지연 읽기 → 힙에 파일 전체가 올라오지 않음
            byLevel = lines.map(BigLogParse::parse).flatMap(Optional::stream)
                .collect(groupingBy(Entry::level, TreeMap::new, counting()));
        }
        try (Stream<String> lines = Files.lines(log)) {
            topErrors = lines.map(BigLogParse::parse).flatMap(Optional::stream)
                .filter(e -> e.level().equals("ERROR"))
                .collect(groupingBy(Entry::message, counting()))
                .entrySet().stream()
                .sorted(Map.Entry.<String, Long>comparingByValue().reversed().thenComparing(Map.Entry.comparingByKey()))
                .limit(2).toList();
        }
        long ms = (System.nanoTime() - t0) / 1_000_000;
        System.out.println("byLevel = " + byLevel);
        System.out.println("topErrors = " + topErrors);
        System.out.println("two passes in " + ms + "ms");
        Files.delete(log);
    }
}
// 출력 (시간 값은 환경에 따라 다름, 건수는 시드 고정이라 동일):
// file size = 12 MB
// byLevel = {ERROR=59956, INFO=180272, WARN=59772}
// topErrors = [db connection refused=20255, cache miss storm=20039]
// two passes in 4796ms

변형 4: 엣지 케이스 — 빈 스트림, null 요소, null 키, toMap 충돌

운영 데이터는 비어 있거나 null 이 섞여 있습니다. 빈 스트림에서 sum 은 0, reduce(op)/max 는 Optional.empty, average 는 OptionalDouble.empty, allMatch 는 공허한 참(true)입니다.

스트림 자체는 null 요소를 허용하지만 groupingBy 는 null 요소·null 키에서, toMap 은 null 값에서 NPE 를 던집니다. 예외 메시지를 보고 어느 층에서 터졌는지 구분하고, filter(Objects::nonNull) 과 requireNonNullElse 로 방어하는 순서를 익힙니다.

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

public class StreamEdgeCases {
    record Order(String id, String customer, Long amount) {}

    public static void main(String[] args) {
        List<Order> none = List.of();

        // 1) 빈 스트림: reduce(identity) 는 항등원, reduce(op)/max 는 empty, average 는 OptionalDouble.empty
        System.out.println(none.stream().mapToLong(Order::amount).sum());
        System.out.println(none.stream().map(Order::amount).reduce(0L, Long::sum));
        System.out.println(none.stream().map(Order::amount).reduce(Long::sum));
        System.out.println(none.stream().max(Comparator.comparing(Order::amount)));
        System.out.println(none.stream().mapToLong(Order::amount).average().orElse(Double.NaN));
        System.out.println(none.stream().collect(groupingBy(Order::customer, counting())));   // 빈 맵, 예외 없음
        System.out.println(none.stream().allMatch(o -> o.amount() > 0) + " " + none.stream().anyMatch(o -> o.amount() > 0));   // 공허한 참

        // 2) null 요소 / null 필드: 스트림 자체는 null 을 허용하지만 컬렉터가 터진다
        List<Order> dirty = Arrays.asList(
            new Order("O1", "kim", 100L), null, new Order("O2", null, 200L), new Order("O3", "lee", null));
        try {
            dirty.stream().collect(groupingBy(Order::customer));
        } catch (NullPointerException e) {
            System.out.println("NPE: null 요소를 filter 하지 않음");
        }
        try {
            dirty.stream().filter(Objects::nonNull).collect(groupingBy(Order::customer));
        } catch (NullPointerException e) {
            System.out.println("NPE: groupingBy 키가 null (element cannot be mapped to a null key)");
        }
        Map<String, Long> clean = dirty.stream()
            .filter(Objects::nonNull)
            .filter(o -> o.customer() != null && o.amount() != null)
            .collect(groupingBy(Order::customer, TreeMap::new, summingLong(Order::amount)));
        System.out.println("clean = " + clean);
        Map<String, Long> withDefaultKey = dirty.stream().filter(Objects::nonNull)
            .collect(groupingBy(o -> Objects.requireNonNullElse(o.customer(), "(unknown)"), TreeMap::new, counting()));
        System.out.println("withDefaultKey = " + withDefaultKey);

        // 3) toMap: null 값은 예외, 중복 키는 merge 로. Collectors.toMap 은 HashMap.merge 를 쓰므로 null 값 불가
        try {
            dirty.stream().filter(Objects::nonNull).collect(toMap(Order::id, Order::amount));
        } catch (NullPointerException e) {
            System.out.println("NPE: toMap 값이 null");
        }
        Map<String, Long> latest = Stream.of(new Order("A", "kim", 1L), new Order("A", "kim", 2L), new Order("B", "lee", 3L))
            .collect(toMap(Order::id, Order::amount, (a, b) -> b, LinkedHashMap::new));   // 마지막 승
        System.out.println("latest = " + latest);
    }
}
// 출력:
// 0
// 0
// Optional.empty
// Optional.empty
// NaN
// {}
// true false
// NPE: null 요소를 filter 하지 않음
// NPE: groupingBy 키가 null (element cannot be mapped to a null key)
// clean = {kim=100}
// withDefaultKey = {(unknown)=1, kim=1, lee=1}
// NPE: toMap 값이 null
// latest = {A=2, B=3}

변형 5: 성능 비교 — 박싱 스트림 vs 기본형 스트림, flatMap vs mapMulti

2.8 의 "박싱 비용"을 2천만 개로 실측합니다. Stream<Integer> 는 map 마다 새 Integer 를 만들고 reduce 에서 언박싱하므로 가장 느리고, mapToInt 로 내려가면 이후 단계가 기본형이 되어 빨라지며, 처음부터 IntStream.range 면 박싱이 아예 없습니다.

1:N 변환에서는 mapMulti 가 요소마다 Stream 객체를 만들지 않아 flatMap 보다 가볍습니다. 네 방식의 합이 같아야(오버플로도 동일하게) 정상입니다.

java
import java.util.*;
import java.util.function.LongSupplier;
import java.util.stream.*;

public class BoxedVsPrimitive {
    static long time(String label, LongSupplier job) {
        long t0 = System.nanoTime();
        long r = job.getAsLong();
        long ms = (System.nanoTime() - t0) / 1_000_000;
        System.out.printf("%-22s = %d  %dms%n", label, r, ms);
        return ms;
    }

    public static void main(String[] args) {
        int n = 20_000_000;
        List<Integer> boxedList = IntStream.range(0, n).boxed().toList();   // Integer 객체 2천만 개

        for (int i = 0; i < 2; i++) {                                         // 워밍업(JIT)
            boxedList.stream().map(x -> x * 2).reduce(0, Integer::sum);
            IntStream.range(0, n).map(x -> x * 2).sum();
        }
        // 같은 계산: 박싱 스트림은 요소마다 Integer 를 새로 만들고 언박싱한다
        time("Stream<Integer> reduce", () -> boxedList.stream().map(x -> x * 2).reduce(0, Integer::sum));
        time("mapToInt + sum", () -> boxedList.stream().mapToInt(Integer::intValue).map(x -> x * 2).sum());
        time("IntStream.range sum", () -> IntStream.range(0, n).map(x -> x * 2).sum());
        time("IntStream parallel", () -> IntStream.range(0, n).parallel().map(x -> x * 2).sum());

        // flatMap vs mapMulti: 1:N 변환에서 mapMulti 는 요소마다 Stream 객체를 만들지 않는다
        List<List<Integer>> nested = IntStream.range(0, 200_000).mapToObj(i -> List.of(i, i, i)).toList();
        for (int i = 0; i < 2; i++) {
            nested.stream().flatMap(List::stream).count();
            nested.stream().<Integer>mapMulti(Iterable::forEach).count();
        }
        time("flatMap count", () -> nested.stream().flatMap(List::stream).count());
        time("mapMulti count", () -> nested.stream().<Integer>mapMulti(Iterable::forEach).count());
    }
}
// 출력 (시간 값은 환경에 따라 다름; int 합은 오버플로로 음수가 될 수 있으나 네 방식이 동일해야 정상):
// Stream<Integer> reduce = 1085788928  1846ms
// mapToInt + sum         = 1085788928  1793ms
// IntStream.range sum    = 1085788928  1263ms
// IntStream parallel     = 1085788928  463ms
// flatMap count          = 600000  241ms
// mapMulti count         = 600000  175ms
응용 변형 예제
  • 변형 1: 같은 집계를 for 문 → 스트림 → 병렬 스트림으로
  • 변형 2: Collector.of 로 만드는 상위 N 컬렉터와 제네릭 topNBy 유틸
  • 변형 3: 30만 줄 로그 파일 스트리밍 파싱
  • 변형 4: 엣지 케이스 — 빈 스트림, null 요소, null 키, toMap 충돌
  • 변형 5: 성능 비교 — 박싱 스트림 vs 기본형 스트림, flatMap vs mapMulti
이전 섹션3 코드 예제4 / 7다음 섹션5 자주 하는 실수 (Tip)