Tất cả bài viết
JavaStream APIFunctional ProgrammingJava 25

Stream API chuyên sâu: từ Collectors nâng cao đến Stream Gatherers

30/09/20260 lượt xem12 phút đọc

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 filter toàn bộ danh sách rồi mới map. 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: findFirst dừ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, takeWhile cũ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ém NullPointerException nếu gặp null.

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. TreeMap mặc định so sánh String theo 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ột Collator.

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.supplyAsync khô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.range chia 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, forEachOrdered buộc phải đồng bộ lại thứ tự. Nếu không cần thứ tự, dùng findAny hoặc unordered().
  • 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 HashSet hay 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ề false khi 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.ofGreedy bá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èm reduce(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ột Supplier<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òng for thườ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.

# Java
# Stream API
# Functional Programming
# Java 25