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

如何确保fs2流订阅后再执行startProducing(类似doOnSubscribe)?

解决fs2 Stream适配非纯Java API的订阅与启动顺序问题

针对你遇到的必须先订阅ImpureProducer的output(Flowable类型)才能调用startProducing的问题,核心是要严格保证订阅动作完成后再执行启动逻辑,彻底消除竞态条件。以下是具体的实现方案:

核心思路

放弃使用parProductL这类并行操作,改为手动控制订阅、启动、消费的顺序:

  1. 用fs2 Queue作为中间层接收Flowable的数据
  2. 手动订阅Flowable,通过Deferred捕捉订阅完成的信号
  3. 确认订阅完成后,再调用startProducing
  4. 从Queue读取数据生成fs2 Stream
  5. 用Resource管理生命周期,确保最终取消订阅并停止producer

具体代码实现

import cats.effect.{IO, Deferred, Resource}
import fs2.{Stream, Queue}
import org.reactivestreams.{Subscription, Flowable}
import java.util.concurrent.atomic.AtomicReference

// 适配非纯Java API的ImpureProducer
class AdaptedProducer(impure: ImpureProducer) {
  def stream: Stream[IO, String] = Stream.resource {
    // 初始化资源:Queue、订阅信号、订阅引用
    for {
      dataQueue <- Queue.unbounded[IO, String]
      subscribedSignal <- Deferred[IO, Unit]
      subscriptionRef <- IO(new AtomicReference[Subscription](null))
      // 手动订阅Flowable,绑定数据到Queue,触发订阅完成信号
      _ <- IO {
        impure.getOutput.subscribe(
          (data: String) => dataQueue.offer(data).unsafeRunSync(),
          (err: Throwable) => IO.raiseError(err).unsafeRunSync(),
          () => dataQueue.close.unsafeRunSync(),
          (s: Subscription) => {
            subscriptionRef.set(s)
            subscribedSignal.complete(()).unsafeRunSync()
          }
        )
      }
      // 等待订阅完成后,启动producer
      _ <- subscribedSignal.get
      _ <- IO(impure.startProducing())
    } yield {
      // 返回消费Queue的流,同时绑定资源释放逻辑
      dataQueue.stream.onFinalize {
        for {
          _ <- IO(Option(subscriptionRef.get()).foreach(_.cancel()))
          _ <- IO(impure.stopProducing())
        } yield ()
      }
    }
  }.flatMap(identity)
}

// 模拟非纯Java API的ImpureProducer
// public class ImpureProducer {
//     private Flowable<String> output = Flowable.just("data1", "data2");
//     private boolean started = false;
//     public Flowable<String> getOutput() { return output; }
//     public void startProducing() { if (started) throw new IllegalStateException(); started = true; }
//     public void stopProducing() { started = false; }
// }

方案优势

  1. 完全消除竞态:订阅动作完成后才会执行startProducing,严格符合API要求
  2. 可靠的数据传递:用fs2 Queue作为中间层,避免Flowable数据丢失
  3. 安全的生命周期管理:通过Resource确保无论流正常结束还是异常终止,都会取消订阅并停止producer
  4. 直观可控:手动订阅逻辑清晰,无需依赖interop库的隐式转换规则

使用示例

val impureProducer = new ImpureProducer()
val adapted = new AdaptedProducer(impureProducer)

// 安全获取第一个输出,无竞态条件
val firstOutput: IO[String] = adapted.stream.head.compile.lastOrError

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 04:36:08