2.3 의 설명을 실행 가능한 코드로 확인합니다. sorted 가 장벽인 것도 함께 봅니다.
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]실무에서 가장 흔한 세 가지 집계입니다. groupingBy + downstream 컬렉터 조합을 봅니다.
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컬렉터의 세부 동작을 한 번에 확인합니다. 특히 toMap 이 키 충돌에서 어떻게 터지고 어떻게 막는지 봅니다.
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]}주문 → 상품 라인을 펼쳐 상품별 판매 수량을 집계하고, 로그 파일을 Files.lines 로 읽어 레벨별 건수와 ERROR 메시지 상위를 뽑습니다.
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]CPU 바운드 작업으로 병렬 이득을 확인하고, 공유 ArrayList 를 쓰면 결과가 깨지는 것과 collect 로 고치는 것을 보여줍니다.
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