如何确保fs2流订阅后再执行startProducing(类似doOnSubscribe)?
解决fs2 Stream适配非纯Java API的订阅与启动顺序问题
针对你遇到的必须先订阅ImpureProducer的output(Flowable类型)才能调用startProducing的问题,核心是要严格保证订阅动作完成后再执行启动逻辑,彻底消除竞态条件。以下是具体的实现方案:
核心思路
放弃使用parProductL这类并行操作,改为手动控制订阅、启动、消费的顺序:
- 用fs2
Queue作为中间层接收Flowable的数据 - 手动订阅Flowable,通过
Deferred捕捉订阅完成的信号 - 确认订阅完成后,再调用
startProducing - 从Queue读取数据生成fs2 Stream
- 用
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; } // }
方案优势
- 完全消除竞态:订阅动作完成后才会执行
startProducing,严格符合API要求 - 可靠的数据传递:用fs2 Queue作为中间层,避免Flowable数据丢失
- 安全的生命周期管理:通过
Resource确保无论流正常结束还是异常终止,都会取消订阅并停止producer - 直观可控:手动订阅逻辑清晰,无需依赖interop库的隐式转换规则
使用示例
val impureProducer = new ImpureProducer() val adapted = new AdaptedProducer(impureProducer) // 安全获取第一个输出,无竞态条件 val firstOutput: IO[String] = adapted.stream.head.compile.lastOrError
内容的提问来源于stack exchange,提问作者MartinHH
相关产品推荐
相关产品推荐

