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

takeWhileInclusive实现是否可安全用于并行序列?

Is this 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 shouldContinue at the same time before it’s updated, leading to multiple threads returning true and processing more elements than intended.
  • Writes to shouldContinue might 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:11:35