Java Programming

Streams, Collectors, Parallel Processing and Date-Time

PGCP-BDA

stream and collection

A collection stores elements, while a stream is a one-use computation pipeline over a source and does not itself own the elements.

stream pipeline

A stream pipeline has a source, lazy intermediate operations and one terminal operation that triggers traversal.

filter and map

filter retains elements satisfying a predicate; map replaces each element with one transformed result.

flatMap

flatMap maps each input to a stream and flattens all resulting streams into one sequence.

reduce

reduce combines stream elements using an associative accumulator and a compatible identity when one is supplied.

collectors

Collectors accumulate stream elements into mutable results such as lists, maps, strings, summaries and grouped structures.

collect performs mutable reduction into containers or summaries:

Map<String, List<Employee>> byDepartment =
        employees.stream().collect(
            Collectors.groupingBy(Employee::department));

Common collectors include:

  • toList, toSet and toCollection;
  • joining for text;
  • groupingBy and partitioningBy;
  • counting, summingInt and averagingDouble;
  • mapping, filtering and collectingAndThen;
  • toMap.

toMap requires a merge function when duplicate keys are possible:

Map<String, Integer> totals = orders.stream().collect(
        Collectors.toMap(
            Order::customer,
            Order::amount,
            Integer::sum));

Without a merge policy, duplicate keys cause IllegalStateException.

grouping and partitioning

groupingBy organizes elements by a classifier key, while partitioningBy always produces Boolean groups from a predicate.

parallel stream

A parallel stream may split work across a common fork-join pool, but order, overhead, blocking, state and data size determine benefit.

associativity

An associative operation can regroup operands without changing the result and is required for deterministic parallel reduction.

LocalDate and Instant

LocalDate represents a date without time zone, while Instant represents a point on the UTC timeline.

Date-time concepts

The java.time API separates distinct meanings:

TypeMeaning
LocalDateCalendar date without time or zone
LocalTimeTime of day without date or zone
LocalDateTimeDate and time without zone
InstantPoint on the UTC time line
OffsetDateTimeDate-time with a fixed UTC offset
ZonedDateTimeDate-time with a region zone and its rules
DurationTime-based amount, such as seconds
PeriodDate-based amount, such as months
LocalDate exam = LocalDate.of(2026, 10, 15);
Instant recorded = Instant.now();
ZonedDateTime india = recorded.atZone(ZoneId.of("Asia/Kolkata"));

Choose the type from meaning. A birthday usually needs LocalDate; an audit event needs an Instant; a scheduled meeting tied to regional daylight-saving rules may need ZonedDateTime.

Encounter order and parallel streams

An ordered source such as List has encounter order; HashSet generally does not promise one. forEach on a parallel stream may process in any order, while forEachOrdered preserves encounter order at a possible performance cost.

parallelStream() is not automatically faster. Splitting, coordination, data size, operation cost, ordering constraints and the shared common ForkJoinPool all matter. Parallel operations should be stateless, non-interfering, associative where required and free from unsafe shared mutation.

// Unsafe idea: several workers mutate the same ArrayList
List<Integer> output = new ArrayList<>();
numbers.parallelStream().forEach(output::add);

Use a collector designed for the reduction instead of mutating shared state.

Laziness and operation fusion

Intermediate operations such as filter, map, peek, distinct and sorted are lazy. Building the pipeline performs no element processing:

Stream<String> pipeline = names.stream()
        .filter(name -> {
            System.out.println("checking " + name);
            return name.startsWith("A");
        });
// No checking yet
long count = pipeline.count();

During traversal, operations can be fused per element. Short-circuiting terminal operations such as findFirst, anyMatch and limit may avoid visiting the whole source. Stateful operations such as sorted or distinct may need to retain information before yielding results.

Choosing loops or streams

Streams are effective for declarative transformations, filtering, grouping and aggregation. A loop may be clearer for stateful algorithms, several exits, checked-exception-heavy processing or logic whose stream version requires hidden mutation.

Readability is the deciding factor. Avoid pipelines so long that the data shape becomes hard to follow. Extract named predicates or functions or divide the transformation into meaningful stages.

Stream model

A stream is a single-use sequence of elements supporting aggregate operations. It is not a data structure and normally does not store elements. A pipeline consists of:

  1. a source, such as a collection, array, generator or file;
  2. zero or more intermediate operations;
  3. one terminal operation.
long passing = List.of(30, 70, 80).stream()
        .filter(mark -> mark >= 40)
        .count();                         // 2

filter is intermediate and count is terminal. The stream pulls values from its source only when the terminal operation begins traversal.

Resource-backed streams

A stream over an in-memory collection normally owns no external resource. Some streams, including those from Files.lines, Files.list and Files.walk, hold open resources and must be closed:

try (Stream<String> lines =
         Files.lines(path, StandardCharsets.UTF_8)) {
    long errors = lines.filter(s -> s.contains("ERROR")).count();
}

Closing a pipeline also invokes registered onClose handlers. Terminal consumption alone should not be assumed to close every resource-backed stream.

Filtering and transformation

filter retains elements for which a Predicate returns true. map replaces each element with one result:

List<String> result = List.of("Neel", "Asha", "Nora").stream()
        .filter(name -> name.startsWith("N"))
        .map(String::toUpperCase)
        .sorted()
        .toList();                         // [NEEL, NORA]

The original list remains unchanged. distinct removes duplicates using equality and hashing. sorted uses natural ordering or a Comparator.

flatMap transforms one input into zero or more output elements and flattens the nested streams:

List<String> words = sentences.stream()
        .flatMap(line -> Arrays.stream(line.split("\\s+")))
        .toList();

Use map for one-to-one transformation and flatMap when each input contains or produces a sequence.

Worked pipeline trace

List<String> result = List.of("Neel", "Asha", "Nora").stream()
        .filter(name -> name.startsWith("N"))
        .map(String::toUpperCase)
        .sorted()
        .toList();

toList triggers traversal. Neel passes and becomes NEEL; Asha is discarded; Nora passes and becomes NORA. Sorting compares the two transformed strings, giving [NEEL, NORA]. The source list is unchanged and the stream is consumed.

Optional fundamentals

Optional<T> represents either one non-null value or absence:

Optional<Student> found = repository.findById(id);
String name = found.map(Student::name)
                   .orElse("Unknown");

Create it with Optional.of(nonNull), ofNullable(possiblyNull) or empty(). of(null) throws NullPointerException.

map transforms a present value and keeps absence. flatMap is used when the mapping function already returns Optional, avoiding Optional<Optional<T>>. filter retains a present value only when it satisfies a predicate.

Terminal operations

Terminal operations produce a result or side effect and consume the stream. Examples include:

  • count, min, max, findFirst and findAny;
  • anyMatch, allMatch and noneMatch;
  • forEach and forEachOrdered;
  • reduce, collect and toList.

A stream cannot ordinarily be used after a terminal operation:

Stream<String> stream = names.stream();
long size = stream.count();
// stream.findFirst(); // IllegalStateException

Create another stream from the reusable source. If stream construction is expensive, a Supplier<Stream<T>> can create a fresh pipeline on demand.

Date-time immutability and arithmetic

Java-time objects are immutable and thread-safe. Arithmetic returns a new value:

LocalDate start = LocalDate.of(2026, 1, 31);
LocalDate later = start.plusMonths(1);

start remains unchanged. Calendar arithmetic adjusts to valid dates, so adding one month to January 31 may produce the last valid day of February. Period.between describes calendar components, while Duration.between describes time-based elapsed units.

Time zones can contain daylight-saving transitions where a local time is skipped or occurs twice. An offset such as +05:30 is fixed; a region ID such as Europe/Paris carries historical and future zone rules.

Reduction

Reduction combines elements into one value:

int total = List.of(2, 3, 4).stream()
        .reduce(0, Integer::sum);           // 9

The identity must be neutral: combining it with any value returns that value. The accumulator and combiner must be associative for regrouping in parallel:

(a operation b) operation c
=
a operation (b operation c)

Addition is associative for mathematical integers, while subtraction is not. Floating-point addition may differ slightly under regrouping because of rounding. A reduction should not mutate shared external state.

The one-argument reduce returns Optional because an empty stream has no value to return.

Parsing and formatting date-time values

DateTimeFormatter is immutable and thread-safe:

DateTimeFormatter formatter =
        DateTimeFormatter.ofPattern("dd MMM uuuu", Locale.ENGLISH);

LocalDate date = LocalDate.parse("26 Sep 2026", formatter);
String text = date.format(formatter);

Parsing text must match the expected formatter. Prefer standard ISO formats for machine exchange. Use explicit patterns and locales for human formats.

Pattern letters are case-sensitive. M represents month, while m represents minute. H is a 24-hour clock, while h is a 12-hour clock generally paired with an AM/PM marker. A plausible-looking output can still be semantically wrong when the pattern is incorrect.

Side effects and debugging

Stream functions should not modify the source or depend on changing external state. Such interference makes results order-dependent and especially unsafe in parallel.

peek observes elements as they flow and is mainly useful for temporary debugging:

long count = values.stream()
        .peek(v -> logger.debug("before: {}", v))
        .filter(this::valid)
        .count();

Because execution is lazy and may be optimized or short-circuited, peek is unsuitable for essential business effects. Use an explicit loop or terminal action when side effects are the purpose.

Primitive streams

IntStream, LongStream and DoubleStream avoid wrapper allocation and add numeric operations:

double average = students.stream()
        .mapToInt(Student::mark)
        .average()
        .orElse(0.0);

mapToInt changes an object stream to IntStream. boxed() converts a primitive stream to a wrapper stream. Range factories distinguish an exclusive upper bound in range(1, 5) from an inclusive upper bound in rangeClosed(1, 5).

Numeric summary statistics can produce count, sum, minimum, maximum and average in one traversal.

Continue learning

Related notes

Put this topic into timed practice

Open mock tests when you want full-exam pacing, or keep drilling in practice mode.