You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Java 8 Stream API连续统计相同元素并生成TokenBag流?

Solution for Consecutive Element Aggregation in Java Stream

Great question! This is a classic run-length encoding scenario—you need to aggregate consecutive identical elements into your TokenBag objects, rather than doing a global group-by (which would lump all "a"s together regardless of their position). Since standard Java Stream operations don't have built-in support for this stateful, sliding aggregation, we'll need a custom approach that maintains Stream's lazy evaluation and doesn't require intermediate collection storage.

Core Approach: Custom Spliterator

The most flexible and idiomatic way to handle this is to create a custom Spliterator that wraps your source Stream's Spliterator. This allows us to track state (current token and count) as we traverse the source elements, emitting a TokenBag whenever we encounter a different element or reach the end of the stream.

Step 1: Your TokenBag Class (for Reference)

First, let's confirm the TokenBag implementation (you already defined this, but it's included here for completeness):

public class TokenBag {
    private String token;
    private int count;

    public TokenBag(String token, int count) {
        this.token = token;
        this.count = count;
    }

    // Getters
    public String getToken() { return token; }
    public int getCount() { return count; }

    @Override
    public String toString() {
        return String.format("(%s, %d)", token, count);
    }
}

Step 2: Custom Spliterator for Consecutive Aggregation

This Spliterator will track the current token and its count, and emit a TokenBag whenever a non-consecutive element is found. It's also designed to support custom consecutive logic for your complex real-world scenario:

import java.util.Spliterator;
import java.util.function.Consumer;

public class ConsecutiveTokenSpliterator implements Spliterator<TokenBag> {
    private final Spliterator<String> sourceSpliterator;
    private String currentToken;
    private int currentCount;
    // Custom predicate to define "consecutive" (flexible for complex logic)
    private final java.util.function.BiPredicate<String, String> isConsecutive;

    // Default constructor (uses equals() for consecutive check)
    public ConsecutiveTokenSpliterator(Spliterator<String> source) {
        this(source, String::equals);
    }

    // Constructor with custom consecutive rules (for your real-world needs)
    public ConsecutiveTokenSpliterator(Spliterator<String> source, java.util.function.BiPredicate<String, String> isConsecutive) {
        this.sourceSpliterator = source;
        this.isConsecutive = isConsecutive;
        this.currentToken = null;
        this.currentCount = 0;
    }

    @Override
    public boolean tryAdvance(Consumer<? super TokenBag> action) {
        try {
            // Traverse source elements until we hit a non-consecutive token or end of stream
            while (sourceSpliterator.tryAdvance(nextToken -> {
                if (currentToken == null) {
                    // Initialize with the first element
                    currentToken = nextToken;
                    currentCount = 1;
                } else if (isConsecutive.test(currentToken, nextToken)) {
                    // Increment count for consecutive token
                    currentCount++;
                } else {
                    // Non-consecutive token: trigger emission of current TokenBag
                    throw new ConsecutiveBreakException(nextToken);
                }
            })) {
                // Loop continues until break or end of source
            }
        } catch (ConsecutiveBreakException e) {
            // Emit the accumulated TokenBag
            action.accept(new TokenBag(currentToken, currentCount));
            // Reset state to the new non-consecutive token
            currentToken = e.getNextToken();
            currentCount = 1;
            return true;
        }

        // Emit the final accumulated TokenBag if we reached end of source
        if (currentToken != null) {
            action.accept(new TokenBag(currentToken, currentCount));
            currentToken = null; // Mark as processed
            return true;
        }

        // No more elements to process
        return false;
    }

    // Helper exception to break out of the inner lambda loop
    private static class ConsecutiveBreakException extends RuntimeException {
        private final String nextToken;

        public ConsecutiveBreakException(String nextToken) {
            this.nextToken = nextToken;
        }

        public String getNextToken() {
            return nextToken;
        }
    }

    @Override
    public Spliterator<TokenBag> trySplit() {
        // Consecutive aggregation is inherently sequential (depends on prior elements)
        // Return null to indicate no parallel splitting support
        return null;
    }

    @Override
    public long estimateSize() {
        // Estimate: worst case (all unique elements) = source size; best case (all same) = 1
        return sourceSpliterator.estimateSize();
    }

    @Override
    public int characteristics() {
        // Inherit source characteristics, remove SIZED (we can't know exact output size)
        return sourceSpliterator.characteristics() & ~Spliterator.SIZED | Spliterator.NONNULL;
    }
}

Step 3: Use the Custom Spliterator to Create Your Stream

Now, wrap your source Stream with the custom Spliterator to get the TokenBag Stream:

import java.util.stream.Stream;
import java.util.stream.StreamSupport;

public class Main {
    public static void main(String[] args) {
        Stream<String> src = Stream.of("a", "a", "a", "b", "b", "a", "a");

        // Create the TokenBag Stream (serial, since parallel doesn't make sense here)
        Stream<TokenBag> tokenBagStream = StreamSupport.stream(
                new ConsecutiveTokenSpliterator(src.spliterator()),
                false
        );

        // Test the output
        tokenBagStream.forEach(System.out::println);
        // Output:
        // (a, 3)
        // (b, 2)
        // (a, 2)
        // Note: Your example mentioned (b, 3) which seems like a typo, since there are only 2 "b"s in the source
    }
}

Generalization for Complex Scenarios

The custom Spliterator is built to handle your real-world requirements:

  • Custom Consecutive Rules: Use the overloaded constructor to pass a BiPredicate<String, String> instead of String::equals. For example, if you need to treat "A" and "a" as consecutive, you could pass (t1, t2) -> t1.equalsIgnoreCase(t2).
  • Lazy Evaluation: The Stream remains lazy—elements are processed only when a terminal operation (like forEach) is called, no intermediate collection is used.
  • Maintain Stream Semantics: You can chain additional intermediate operations (like filter, map) to the TokenBag Stream just like any other Java Stream.

Why Not Other Approaches?

  • Collectors.groupingBy: This does global aggregation, not consecutive grouping—you'd get one TokenBag for "a" (count 5) and one for "b" (count 2), which isn't what you want.
  • Stream.reduce: This is a terminal operation that returns a single result, not a Stream of results.
  • Third-Party Libraries: While libraries like RxJava or Vavr have built-in operators for this, sticking to standard Java gives you better portability.

内容的提问来源于stack exchange,提问作者Nick Legend

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:24:30