如何用Java 8 Stream API连续统计相同元素并生成TokenBag流?
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 ofString::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 theTokenBagStream just like any other Java Stream.
Why Not Other Approaches?
Collectors.groupingBy: This does global aggregation, not consecutive grouping—you'd get oneTokenBagfor "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

