Most Java stream pipelines end in collect(...), and most of the interesting work in them happens there: grouping orders by customer, building a lookup map, computing several statistics in one pass. Collectors is a library of ready-made recipes for that terminal step, and it is easy to use by pattern-matching on examples. It is also where a surprising share of stream bugs live: a toMap that throws in production on the first duplicate key, a groupingBy that dies on a null, a parallel collect that is slower than the sequential one.

This article explains collectors from the contract up. Once you know the four functions and three flags a Collector is made of, every factory method in Collectors becomes predictable, including how it behaves in parallel and when it fails. The intermediate side of pipelines is covered in Java Stream operations; here we stay on the terminal side and finish with a custom collector and a checklist for reviewing collect calls in your own code.

Advertisement

The contract: four functions and a set of flags

A Collector<T, A, R> describes a mutable reduction. T is the element type, A is the mutable accumulation container, usually hidden from callers, and R is the result. It is defined by five methods:

MethodTypeJob
supplier()Supplier<A>Create a new, empty container.
accumulator()BiConsumer<A, T>Fold one element into a container, mutating it.
combiner()BinaryOperator<A>Merge two partial containers; used when a parallel stream was split.
finisher()Function<A, R>Turn the final container into the result.
characteristics()Set<Characteristics>Hints: IDENTITY_FINISH, UNORDERED, CONCURRENT.

A sequential collect is simple: call the supplier once, call the accumulator for every element in encounter order, then call the finisher. If the collector reports IDENTITY_FINISH, the container already is the result and the finisher is skipped with an unchecked cast.

The contract imposes two rules on implementers. Identity: combining a container with an empty one must give the same result as the container alone. Associativity: splitting the input anywhere and combining the pieces must give the same result as accumulating everything into one container. If your collector breaks either rule it will look correct in sequential tests and give different answers in parallel, which is the worst kind of bug to find.

How a parallel collect really runs

On a parallel stream the source spliterator is split into chunks that run as fork/join tasks. Each leaf task calls the supplier to get its own container and accumulates its chunk into it, so no container is ever touched by two threads. As tasks complete, the framework calls the combiner to merge sibling results, left with right, up the tree, and calls the finisher once on the root. Because merging keeps left before right, encounter order is preserved even though the chunks ran concurrently.

A parallel collect: one container per chunk, merged by the combiner, finished onceSource spliteratorsplit into chunkssupplier()container A1supplier()container A2supplier()container A3supplier()container A4accumulator(A, t) folds each element of a chunk into that chunk's own container (no sharing, no locks)combiner(A1, A2)combiner(A3, A4)combiner(L, R)finisher(A)skipped if IDENTITY_FINISHCONCURRENT + UNORDERED on a parallel stream: one shared container, no combiner calls
Each leaf task owns a container; the combiner merges pairs up the split tree; the finisher runs once. A concurrent collector replaces the tree with one shared, thread-safe container.

That design is why ordinary collectors are thread-safe without locks, and it also explains their cost. Merging two ArrayLists copies the right one into the left. Merging two HashMaps in groupingBy re-inserts every key of one map into the other and merges the per-key downstream containers. For a grouping with many keys, the merge phase can cost as much as the accumulation, which is one reason parallel groupingBy is often slower than sequential. The general rules for when parallelism pays are in Java parallel streams.

The alternative is a concurrent reduction: one shared container, such as a ConcurrentHashMap, that all threads accumulate into directly, with no combiner calls. The stream only does this when three conditions hold together: the stream is parallel, the collector has the CONCURRENT characteristic, and either the stream is unordered or the collector has UNORDERED. groupingByConcurrent and toConcurrentMap carry both flags. The price is ordering: the lists inside a concurrent grouping hold elements in whatever order threads arrived.

Advertisement

Collecting into lists and sets: three versions with different guarantees

There are three common ways to get a List out of a stream and they are not interchangeable:

CallSinceMutabilityNulls
collect(Collectors.toList())8Unspecified; the javadoc makes no guarantee about type or mutabilityAllowed
collect(Collectors.toUnmodifiableList())10UnmodifiableNullPointerException on a null element
stream.toList()16UnmodifiableAllowed

Code that adds to the result of Collectors.toList() works today only because the current implementation happens to return a mutable list. If you need a mutable list, say so with Collectors.toCollection(ArrayList::new). If you do not, prefer stream.toList() for its brevity and its unmodifiable result. The same three-way split exists for sets with toSet, toUnmodifiableSet and toCollection(TreeSet::new), the last being the way to get a sorted or insertion-ordered set.

toMap: the two exceptions everyone meets

The two-argument toMap(keyMapper, valueMapper) assumes keys are unique. On the first duplicate it throws IllegalStateException with a message naming the key and both values. Collecting users by email works in every test fixture and fails the day two accounts share an address in different letter case. Decide what a duplicate means and pass a merge function:

// Keep the most recently updated user per normalised email.
Map<String, User> byEmail = users.stream()
    .collect(Collectors.toMap(
        u -> u.email().toLowerCase(Locale.ROOT),
        Function.identity(),
        (a, b) -> a.updatedAt().isAfter(b.updatedAt()) ? a : b,
        LinkedHashMap::new));            // fourth argument picks the map type

The second exception is less well known: toMap throws NullPointerException if the value mapper returns null. Every overload rejects null values, even though HashMap itself accepts them; the forms with a merge function do it through Map.merge. Either filter the nulls out, map them to an explicit sentinel or Optional, or fall back to a plain loop when null values are meaningful. Null keys behave differently across the overloads, so the safe rule is not to produce them at all. When a value is legitimately absent, an explicit Optional is clearer; see Java Optional.

groupingBy and downstream collectors

groupingBy(classifier) returns a Map<K, List<T>> built in a HashMap. Its real power is the downstream argument: a second collector applied to each group. Collectors compose, so a single pass over the data can produce nested summaries.

record Order(String customer, String region, String status, long cents, List<String> skus) {}

// Revenue per region, sorted by region name.
Map<String, Long> revenue = orders.stream()
    .collect(Collectors.groupingBy(Order::region, TreeMap::new,
             Collectors.summingLong(Order::cents)));

// Per customer: number of orders, and the distinct SKUs bought, only for PAID orders.
Map<String, Set<String>> skusByCustomer = orders.stream()
    .collect(Collectors.groupingBy(Order::customer,
             Collectors.filtering(o -> o.status().equals("PAID"),
             Collectors.flatMapping(o -> o.skus().stream(), Collectors.toSet()))));

// Largest order per region, unwrapped from Optional.
Map<String, Order> biggest = orders.stream()
    .collect(Collectors.groupingBy(Order::region,
             Collectors.collectingAndThen(
                 Collectors.maxBy(Comparator.comparingLong(Order::cents)),
                 Optional::get)));

Some details to know. counting() produces a Long and the averaging collectors produce a Double that is 0.0 for an empty input. filtering differs from a filter before the grouping: filtering inside keeps a key whose elements were all removed, mapped to an empty group, while filtering before drops the key altogether. That matters for reports that should list every region. maxBy returns an Optional because the collector cannot know a group is non-empty; inside groupingBy it never is empty, which is why Optional::get is safe there and nowhere else.

The classifier must not return null; groupingBy throws NullPointerException with the message that an element cannot be mapped to a null key. Map missing values to a named bucket such as "UNKNOWN" before grouping. partitioningBy(predicate) is the two-way special case, and unlike groupingBy its result always contains both true and false keys, even if one side is empty.

teeing, joining and summary statistics

teeing(c1, c2, merger), added in JDK 12, sends every element to two collectors and merges their results. It replaces the habit of streaming the same list twice:

record Range(long min, long max) {}

Range r = orders.stream().map(Order::cents)
    .collect(Collectors.teeing(
        Collectors.minBy(Comparator.naturalOrder()),
        Collectors.maxBy(Comparator.naturalOrder()),
        (min, max) -> new Range(min.orElse(0L), max.orElse(0L))));

// For numeric summaries a single statistics object is simpler still:
LongSummaryStatistics s = orders.stream().collect(Collectors.summarizingLong(Order::cents));
// s.getCount(), s.getSum(), s.getMin(), s.getMax(), s.getAverage()

joining(delimiter, prefix, suffix) builds a string with a StringBuilder and is the right tool for producing CSV lines or log messages from a stream of strings. reducing exists mainly for use as a downstream; at the top level, Stream.reduce or a primitive stream's sum is clearer and avoids boxing.

Writing your own collector

When no combination of factory methods fits, write a collector with Collector.of. A realistic example: the top N elements by some score, computed in a single pass with bounded memory. The container is a min-heap of size N; the combiner merges two heaps and trims back to N, which keeps both identity and associativity.

static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> cmp) {
    return Collector.of(
        () -> new PriorityQueue<T>(cmp),                // min-heap: smallest kept element on top
        (heap, t) -> {
            heap.offer(t);
            if (heap.size() > n) heap.poll();          // drop the smallest
        },
        (left, right) -> {
            for (T t : right) {
                left.offer(t);
                if (left.size() > n) left.poll();
            }
            return left;
        },
        heap -> {
            List<T> out = new ArrayList<>(heap);
            out.sort(cmp.reversed());                  // largest first
            return out;
        });
}

List<Order> top5 = orders.parallelStream()
    .collect(topN(5, Comparator.comparingLong(Order::cents)));

Test a custom collector the way the framework will use it. Run it sequentially, then split a list by hand into two halves, accumulate each into its own container, combine, finish and assert the same result. Do it with an empty half too, to check identity. Ten lines of test catch the collectors whose combiner returns the wrong container or forgets to apply the same invariant as the accumulator. Do not declare CONCURRENT unless your container really is safe for concurrent accumulation, and do not declare IDENTITY_FINISH if the finisher does anything, since it will be skipped.

Worked example: one pass for a sales dashboard

Suppose a nightly job loads two million orders and must produce, per region, the revenue, the order count and the five biggest orders. A first draft streams the list three times with three separate groupings. Using the pieces above it becomes one pass:

record RegionStats(long revenue, long count, List<Order> top) {}

Map<String, RegionStats> dashboard = orders.stream()
    .filter(o -> o.region() != null)
    .collect(Collectors.groupingBy(Order::region, TreeMap::new,
        Collectors.teeing(
            Collectors.summarizingLong(Order::cents),
            topN(5, Comparator.comparingLong(Order::cents)),
            (stats, top) -> new RegionStats(stats.getSum(), stats.getCount(), top))));

The data flow is: each order is classified by region; inside the region's group, teeing feeds it to both a statistics accumulator and the top-N heap; at the end each group's finisher builds a RegionStats. Memory is proportional to the number of regions times five orders, not to the input. If profiling later shows the job is CPU-bound and the region count is small, switching to parallelStream() is safe because every piece honours the contract, and the merge cost is low because there are few keys. If the keys were customer IDs in the millions, the map merges would dominate and the sequential version would likely stay faster; measure with JMH before deciding.

Failure modes and how they look in production

  • Duplicate key in toMap. IllegalStateException: Duplicate key from a job that passed every test. Fix: decide the merge rule explicitly; never rely on data being unique.
  • Null value in toMap or null classifier in groupingBy. NullPointerException from deep inside Collectors. Fix: map absent values to a sentinel or filter first.
  • Mutating an unmodifiable result. UnsupportedOperationException after switching from collect(toList()) to toList(). Fix: use toCollection where the caller mutates.
  • Side effects in the accumulator. Writing to an outer ArrayList from forEach or from a lambda inside a collector loses elements in parallel. Fix: let the collector own all state.
  • Slow parallel grouping. High-cardinality groupingBy on a parallel stream spends its time merging maps. Fix: stay sequential, or use groupingByConcurrent when order inside groups does not matter.
  • Hidden boxing. summingInt over a Stream<Integer> unboxes every element. For a single total, mapToInt(...).sum() is faster and clearer.

When not to use a collector

Collectors are a good fit for building containers and summaries. They are a poor fit when you need early exit (use anyMatch or findFirst), when the logic has several branches that update several structures with checked exceptions, or when the result is a single primitive (use a primitive stream). A plain loop is not a failure of style; it is the better choice when the collector expression needs a comment to explain what it builds. For the underlying container types and their costs, see the Java Collections Framework, and for the rest of the stream API, the Java Streams API.

What to do next

  1. Search your codebase for Collectors.toMap( calls with two arguments and add an explicit merge function or a comment proving keys are unique.
  2. Find value mappers that can return null in toMap and classifiers that can return null in groupingBy, and map them to sentinels.
  3. Replace collect(Collectors.toList()) with toList() where callers do not mutate, and with toCollection(ArrayList::new) where they do.
  4. Look for pipelines that stream the same source two or three times and fold them into one pass with teeing or a downstream collector.
  5. For any custom collector, add a split-and-combine unit test including an empty half.
  6. Before making a grouping parallel, benchmark it with JMH at production key cardinality.
Key takeaway: A collector is four functions and a set of flags: supplier, accumulator, combiner, finisher and characteristics. Sequential collects use one container; parallel collects give every chunk its own container and merge them with the combiner, which is safe without locks but can be expensive for large maps. Concurrent reduction needs a parallel stream, a CONCURRENT collector and unordered semantics. Most production failures come from toMap duplicate keys, null values and null grouping keys, or mutating an unmodifiable list. Compose downstream collectors and teeing to compute summaries in one pass, and test any custom collector by splitting and combining it by hand.