Stream API chuyên sâu: từ Collectors nâng cao đến Stream Gatherers
Stream thực sự chạy như thế nào?
Hầu hết lập trình viên Java dùng Stream hằng ngày, nhưng ít người hình dung chính xác thứ tự thực thi của nó. Một stream pipeline gồm ba phần: source (nguồn dữ liệu), các intermediate operation (filter, map, ...) và một terminal operation (toList, findFirst, collect, ...). Intermediate operation là lazy: chúng chỉ mô tả việc cần làm, không chạy gì cả cho đến khi gặp terminal operation.
String result = Stream.of("a", "bb", "ccc", "dddd")
.filter(s -> {
System.out.println("filter: " + s);
return s.length() > 1;
})
.map(s -> {
System.out.println("map: " + s);
return s.toUpperCase();
})
.findFirst()
.orElseThrow();
System.out.println("Kết quả: " + result);
filter: a
filter: bb
map: bb
Kết quả: BB
Hai điều đáng chú ý:
- Xử lý theo chiều dọc: Stream không
filtertoàn bộ danh sách rồi mớimap. Từng phần tử đi xuyên qua cả pipeline trước khi phần tử tiếp theo bắt đầu. - Short-circuit:
findFirstdừng ngay khi có kết quả, nên"ccc"và"dddd"không bao giờ được xử lý. Các operation nhưanyMatch,limit,takeWhilecũng vậy.
Ngoại lệ là các stateful operation như sorted: để sắp xếp, nó buộc phải đọc hết mọi phần tử trước khi đẩy phần tử đầu tiên xuống dưới. Đặt sorted() trước một filter loại bỏ 99% dữ liệu là cách lãng phí phổ biến.
toList() và Collectors.toList() không giống nhau
List<Integer> a = Stream.of(1, 2).toList(); // Java 16+
List<Integer> b = Stream.of(1, 2).collect(Collectors.toList());
b.getClass().getSimpleName(); // ArrayList
a.add(3); // UnsupportedOperationException
Stream.of("x", null).toList(); // [x, null]
Stream.toList()trả về list không sửa được và cho phép phần tửnull.Collectors.toList()không cam kết gì về kiểu hay khả năng sửa đổi. Hiện tại nó trả vềArrayList, nhưng code dựa vào điều đó là dựa vào chi tiết cài đặt.Collectors.toUnmodifiableList()cũng không sửa được, nhưng némNullPointerExceptionnếu gặpnull.
mapMulti: khi flatMap là quá nặng
flatMap tạo ra một Stream mới cho mỗi phần tử. Nếu mỗi phần tử chỉ sinh 0–2 kết quả, mapMulti (Java 16) gọn và rẻ hơn: bạn đẩy thẳng kết quả vào một consumer.
List<Integer> expanded = Stream.of(1, 2, 3)
.<Integer>mapMulti((n, downstream) -> {
if (n % 2 == 1) {
downstream.accept(n);
downstream.accept(n * 10);
}
})
.toList();
// [1, 10, 3, 30]
Collectors nâng cao
Các ví dụ dưới đây dùng chung một tập dữ liệu đơn hàng:
record Order(String customer, String city, String category, int amount) {}
List<Order> orders = List.of(
new Order("An", "Hà Nội", "Sách", 120),
new Order("Bình", "HCM", "Điện tử", 900),
new Order("An", "Hà Nội", "Điện tử", 650),
new Order("Chi", "Đà Nẵng", "Sách", 80),
new Order("Dũng", "HCM", "Sách", 200),
new Order("Bình", "HCM", "Thời trang", 300)
);
groupingBy với downstream collector
groupingBy có ba tham số: hàm lấy key, factory tạo Map, và một downstream collector quyết định làm gì với các phần tử trong từng nhóm. Đây là chỗ sức mạnh thật sự nằm.
Map<String, Long> countByCity = orders.stream()
.collect(Collectors.groupingBy(Order::city, TreeMap::new, Collectors.counting()));
// {HCM=3, Hà Nội=2, Đà Nẵng=1}
Map<String, Integer> revenueByCity = orders.stream()
.collect(Collectors.groupingBy(Order::city, TreeMap::new, Collectors.summingInt(Order::amount)));
// {HCM=1400, Hà Nội=770, Đà Nẵng=80}
Map<String, Set<String>> customersByCategory = orders.stream()
.collect(Collectors.groupingBy(
Order::category,
TreeMap::new,
Collectors.mapping(Order::customer, Collectors.toCollection(TreeSet::new))
));
// {Sách=[An, Chi, Dũng], Thời trang=[Bình], Điện tử=[An, Bình]}
Để ý thứ tự: "HCM" đứng trước "Hà Nội", "Đà Nẵng" và "Điện tử" bị đẩy xuống cuối.
TreeMapmặc định so sánhStringtheo mã Unicode: chữ hoa đứng trước chữ thường có dấu, còn "Đ" (U+0110) đứng sau toàn bộ bảng chữ cái Latin cơ bản. Muốn sắp xếp đúng kiểu tiếng Việt, hãy truyền mộtCollator.
List<String> cities = List.of("Đà Nẵng", "HCM", "Hà Nội", "Cần Thơ", "Huế");
new TreeSet<>(cities);
// [Cần Thơ, HCM, Huế, Hà Nội, Đà Nẵng]
TreeSet<String> vi = new TreeSet<>(Collator.getInstance(Locale.of("vi")));
vi.addAll(cities);
// [Cần Thơ, Đà Nẵng, Hà Nội, HCM, Huế]
partitioningBy
Trường hợp đặc biệt của groupingBy khi key là boolean. Kết quả luôn có đủ cả hai key true và false, kể cả khi một nhóm rỗng.
Map<Boolean, List<String>> bigOrders = orders.stream()
.collect(Collectors.partitioningBy(
o -> o.amount() >= 500,
Collectors.mapping(Order::customer, Collectors.toList())
));
// {false=[An, Chi, Dũng, Bình], true=[Bình, An]}
toMap và hai cái bẫy
Bẫy thứ nhất: key trùng. Khách hàng "An" có hai đơn hàng, và toMap không biết phải giữ giá trị nào:
orders.stream().collect(Collectors.toMap(Order::customer, Order::amount));
java.lang.IllegalStateException: Duplicate key An (attempted merging values 120 and 650)
Cách sửa là truyền merge function — ở đây là cộng dồn:
Map<String, Integer> totalByCustomer = orders.stream()
.collect(Collectors.toMap(Order::customer, Order::amount, Integer::sum, TreeMap::new));
// {An=770, Bình=1200, Chi=80, Dũng=200}
Bẫy thứ hai: giá trị null. Collectors.toMap ném NullPointerException nếu hàm lấy value trả về null, kể cả khi HashMap vốn cho phép value null. Nếu dữ liệu có thể thiếu, hãy lọc trước hoặc thay null bằng giá trị mặc định.
teeing: hai collector trong một lần duyệt
Collectors.teeing (Java 12) đưa mỗi phần tử vào hai collector cùng lúc rồi gộp kết quả. Dùng khi cần nhiều thống kê mà không muốn duyệt dữ liệu nhiều lần:
record Summary(long orderCount, int revenue) {}
Summary summary = orders.stream()
.collect(Collectors.teeing(
Collectors.counting(),
Collectors.summingInt(Order::amount),
Summary::new
));
// Summary[orderCount=6, revenue=2250]
Parallel streams: khi nào nhanh, khi nào hại
Thêm .parallel() trông như một cách tăng tốc miễn phí. Thực tế, nó dễ làm hại nhiều hơn là giúp. Xem đoạn code sau:
List<Integer> result = new ArrayList<>();
IntStream.range(0, 10_000).parallel().forEach(result::add); // ❌ data race
Chạy ba lần trên máy mình, kết quả result.size() lần lượt là 10000, 4235 và 2965. Không exception, không cảnh báo — dữ liệu mất một cách âm thầm, và lần chạy đầu tiên còn cho kết quả đúng, đủ để test pass. ArrayList không thread-safe, và forEach song song ghi vào nó từ nhiều thread. Cách đúng là để stream tự gom kết quả: IntStream.range(0, 10_000).parallel().boxed().toList().
Những điều cần biết trước khi dùng parallel():
- Nó chạy trên
ForkJoinPool.commonPool(), một pool dùng chung cho cả JVM (CompletableFuture.supplyAsynckhông truyền executor cũng dùng pool này). Gọi API hay query database bên trong parallel stream sẽ chiếm hết thread của pool và làm chậm những phần khác của ứng dụng. - Nguồn dữ liệu phải chia nhỏ được hiệu quả:
ArrayList, mảng,IntStream.rangechia rất tốt;LinkedList,Stream.iterate,BufferedReader.lines()chia rất kém. - Operation có thứ tự tốn kém khi chạy song song:
findFirst,limit,forEachOrderedbuộc phải đồng bộ lại thứ tự. Nếu không cần thứ tự, dùngfindAnyhoặcunordered(). - Chỉ đáng khi khối lượng tính toán lớn: số phần tử nhân với chi phí xử lý mỗi phần tử phải đủ lớn để bù chi phí chia việc và gộp kết quả. Với vài nghìn phần tử và phép tính nhẹ, bản tuần tự thường nhanh hơn. Luôn đo bằng JMH trước khi quyết định.
Stream Gatherers (Java 24+)
Suốt 10 năm, Stream API có một giới hạn khó chịu: bạn có thể viết terminal operation tùy biến (qua Collector), nhưng không thể viết intermediate operation tùy biến. Muốn chia stream thành từng nhóm 3 phần tử, tính trung bình trượt, hay loại trùng theo một thuộc tính — bạn phải thoát khỏi stream, dùng vòng lặp, hoặc nhờ thư viện ngoài.
Stream Gatherers (JEP 485, chính thức từ Java 24) lấp đúng khoảng trống này với method mới Stream.gather(Gatherer). Có thể xem Gatherer là "Collector cho intermediate operation".
Các gatherer có sẵn
Class java.util.stream.Gatherers cung cấp năm gatherer dùng ngay:
Stream.of(1, 2, 3, 4, 5, 6, 7, 8).gather(Gatherers.windowFixed(3)).toList();
// [[1, 2, 3], [4, 5, 6], [7, 8]]
Stream.of(1, 2, 3, 4, 5).gather(Gatherers.windowSliding(3)).toList();
// [[1, 2, 3], [2, 3, 4], [3, 4, 5]]
Stream.of(1, 2, 3, 4).gather(Gatherers.scan(() -> 0, Integer::sum)).toList();
// [1, 3, 6, 10] (tổng lũy kế)
Stream.of("J", "a", "v", "a").gather(Gatherers.fold(() -> "", String::concat)).toList();
// [Java] (gộp thành một phần tử duy nhất)
Ví dụ thực tế với windowSliding: trung bình trượt 3 phiên của giá cổ phiếu.
List<Double> prices = List.of(10.0, 11.0, 12.0, 13.0, 14.0);
List<Double> movingAvg = prices.stream()
.gather(Gatherers.windowSliding(3))
.map(w -> w.stream().mapToDouble(Double::doubleValue).average().orElseThrow())
.toList();
// [11.0, 12.0, 13.0]
mapConcurrent: gọi I/O song song mà vẫn giữ thứ tự
Gatherers.mapConcurrent(maxConcurrency, mapper) chạy mapper trên virtual threads, tối đa maxConcurrency tác vụ cùng lúc, và trả kết quả theo đúng thứ tự đầu vào. Đây chính xác là thứ bạn muốn khi gọi một API cho danh sách ID:
List<String> results = IntStream.rangeClosed(1, 10).boxed()
.gather(Gatherers.mapConcurrent(5, id -> {
Thread.sleep(200); // giả lập gọi API mất 200ms
return "user-" + id;
}))
.toList();
// [user-1, user-2, ..., user-10] sau khoảng 400ms
Mười lời gọi mỗi cái 200ms, chạy tuần tự mất 2 giây; với 5 luồng đồng thời chỉ còn khoảng 400ms. Tham số maxConcurrency đồng thời là cơ chế backpressure: bạn không vô tình bắn 10.000 request cùng lúc vào một service khác. So với parallel(), cách này không đụng tới common pool và không bị giới hạn bởi số core. (Đoạn code trên được rút gọn: trong code thật, Thread.sleep ném checked exception nên cần bọc try/catch.)
Cấu trúc của một Gatherer
Một Gatherer<T, A, R> nhận phần tử kiểu T, giữ trạng thái kiểu A, và đẩy xuống phần tử kiểu R. Nó gồm bốn thành phần, giống cấu trúc của Collector:
- initializer: tạo trạng thái ban đầu (ví dụ một
HashSethay một buffer). - integrator: được gọi với mỗi phần tử; có thể đẩy 0, 1 hoặc nhiều phần tử xuống downstream. Trả về
falseđể dừng sớm. - combiner: gộp trạng thái khi chạy song song. Bỏ qua nếu gatherer chỉ chạy tuần tự.
- finisher: chạy khi hết đầu vào, dùng để "xả" phần còn lại trong trạng thái.
Tự viết gatherer: distinctBy
distinct() chỉ loại trùng dựa trên equals của cả object. Rất thường ta muốn loại trùng theo một thuộc tính, ví dụ email:
static <T, K> Gatherer<T, ?, T> distinctBy(Function<? super T, ? extends K> keyExtractor) {
return Gatherer.ofSequential(
HashSet<K>::new, // initializer
Gatherer.Integrator.ofGreedy((seen, element, downstream) -> {
if (seen.add(keyExtractor.apply(element))) {
return downstream.push(element); // lần đầu gặp key: đẩy xuống
}
return true; // trùng: bỏ qua, tiếp tục
})
);
}
record User(String email, String name) {}
List<User> users = List.of(
new User("an@example.com", "An"),
new User("binh@example.com", "Bình"),
new User("an@example.com", "An (trùng)")
);
users.stream().gather(distinctBy(User::email)).map(User::name).toList();
// [An, Bình]
Hai chi tiết quan trọng:
downstream.push(element)trả vềfalsekhi downstream không muốn nhận thêm (ví dụ phía sau cólimit(10)và đã đủ 10). Trả tiếp giá trị đó ra ngoài giúp cả pipeline dừng sớm.Integrator.ofGreedybáo cho runtime biết integrator này không tự dừng sớm, cho phép tối ưu tốt hơn.
Gatherer có finisher: batch
Ví dụ sau cần finisher: gom phần tử thành từng lô, và khi hết đầu vào thì đẩy nốt lô cuối dù chưa đủ kích thước.
static <T> Gatherer<T, ?, List<T>> batch(int size) {
return Gatherer.ofSequential(
() -> new ArrayList<T>(size),
Gatherer.Integrator.ofGreedy((buffer, element, downstream) -> {
buffer.add(element);
if (buffer.size() == size) {
List<T> full = List.copyOf(buffer);
buffer.clear();
return downstream.push(full);
}
return true;
}),
(buffer, downstream) -> { // finisher
if (!buffer.isEmpty()) {
downstream.push(List.copyOf(buffer));
}
}
);
}
IntStream.rangeClosed(1, 7).boxed().gather(batch(3)).toList();
// [[1, 2, 3], [4, 5, 6], [7]]
Trên thực tế Gatherers.windowFixed đã làm đúng việc này, nhưng batch là khuôn mẫu tốt cho các biến thể phức tạp hơn — ví dụ gom theo tổng dung lượng thay vì số phần tử. Ứng dụng quen thuộc: insert vào database theo lô 500 dòng thay vì từng dòng một.
Kết hợp gatherer
Gatherer ghép được với nhau bằng andThen, tạo thành một gatherer mới có thể tái sử dụng:
Stream.of("a", "b", "a", "c", "d", "b", "e")
.gather(distinctBy((String s) -> s).andThen(Gatherers.windowFixed(2)))
.toList();
// [[a, b], [c, d], [e]]
Vài nguyên tắc về hiệu năng
- Tránh boxing: với số, dùng
IntStream/LongStream(mapToInt(...).sum()) thay vìStream<Integer>kèmreduce(0, Integer::sum). - Stream chỉ dùng được một lần: gọi terminal operation lần thứ hai trên cùng một stream sẽ ném
IllegalStateException: stream has already been operated upon or closed. Nếu cần duyệt nhiều lần, giữ lại nguồn dữ liệu hoặc mộtSupplier<Stream<T>>. - Đặt filter càng sớm càng tốt, và các stateful operation (
sorted,distinct) càng muộn càng tốt. - Stream không phải lúc nào cũng tốt hơn vòng lặp: với logic có nhiều nhánh, biến trạng thái và
break, một vòngforthường dễ đọc hơn. Stream mạnh nhất khi bài toán thực sự là một chuỗi biến đổi dữ liệu.
Kết luận
Hiểu cách Stream thực thi — lazy, theo chiều dọc, short-circuit — giúp bạn viết pipeline vừa đúng vừa nhanh. Collectors như groupingBy với downstream, toMap với merge function và teeing giải quyết phần lớn bài toán tổng hợp dữ liệu mà không cần vòng lặp thủ công. Và với Stream Gatherers từ Java 24, mảnh ghép cuối cùng đã có: giờ bạn có thể tự định nghĩa intermediate operation, từ cửa sổ trượt đến gọi API song song có kiểm soát, mà vẫn giữ nguyên phong cách khai báo của Stream API.