takeWhileInclusive实现是否可安全用于并行序列?
takeWhileInclusive implementation safe for parallel sequences? Great question—you’re absolutely right to be concerned about thread safety here. This implementation is not safe for parallel sequences, and your worry about the unsynchronized shouldContinue variable is spot-on. Let’s break down why, and how to fix it.
Why the original code is unsafe
The core issue is that shouldContinue is a mutable shared variable accessed and modified by multiple threads when using a parallel sequence. Without synchronization:
- Multiple threads could read
shouldContinueat the same time before it’s updated, leading to multiple threads returningtrueand processing more elements than intended. - Writes to
shouldContinuemight not be visible to other threads immediately (due to CPU caching), causing some threads to keep processing elements even after the predicate has failed for another element. - The order of reads and writes could be reordered by the JVM’s optimizations, leading to unpredictable behavior like early termination or including extra elements.
The original logic works perfectly for sequential sequences because there’s only one thread executing the takeWhile lambda, so there’s no race condition. But parallel execution breaks this guarantee.
A thread-safe alternative
To fix this, we need to replace the plain var with a thread-safe atomic variable that ensures reads and writes are atomic and visible across threads. Here’s an updated implementation using AtomicBoolean:
import java.util.concurrent.atomic.AtomicBoolean fun <T> Sequence<T>.takeWhileInclusive(pred: (T) -> Boolean): Sequence<T> { val shouldContinue = AtomicBoolean(true) return takeWhile { element -> val currentState = shouldContinue.get() if (currentState) { val matchesPredicate = pred(element) if (!matchesPredicate) { shouldContinue.set(false) } } currentState } }
AtomicBoolean provides atomic get() and set() operations, so there’s no way for two threads to interfere with each other’s reads/writes to the state variable. This ensures that once shouldContinue is set to false, all subsequent threads will see that value and stop processing elements.
A caveat about parallel sequence semantics
Even with this thread-safe fix, you should be aware that takeWhileInclusive has ambiguous semantics in parallel sequences. By nature, parallel sequences process elements out of order—so the "first element that fails the predicate" might not be the one you expect in a sequential context. If your logic depends on processing elements in a specific order to determine when to stop, a parallel sequence might not be the right choice here.
内容的提问来源于stack exchange,提问作者jivimberg

