Akka流图中如何复制Source[ByteString, Any]到双Sink并获取Source类型结果?
问题分析
你当前代码的问题在于使用了Sink.last,它的作用是收集流的最后一个元素并返回Future[ByteString],而非你需要的Source[ByteString, Any]。要实现输入流的复制(分流到两个目的地,一个转成InputStream,另一个保留为Source),需要使用Akka流的广播机制来让输入流的元素同时流向两个分支。
解决方案
可以借助BroadcastHub来创建一个可共享的Source,输入流的元素会被广播到所有订阅这个Source的消费者。具体修改如下:
private def duplicateStream(content: Source[ByteString, Any]): Future[Either[X, Y]] = { // 创建BroadcastHub,将输入流转换为可广播的Source val sharedSource: Source[ByteString, NotUsed] = content.runWith(BroadcastHub.sink(bufferSize = 256)) // 第一个分支:将sharedSource转换为InputStream val isFuture: Future[InputStream] = sharedSource.runWith(StreamConverters.asInputStream()) // 第二个分支:直接使用sharedSource作为你需要的bs val bs: Source[ByteString, NotUsed] = sharedSource // 这里根据业务逻辑处理isFuture和bs,最终返回Future[Either[X,Y]] isFuture.map { is => // 处理InputStream的业务逻辑 Right(...) // 替换为你的Y类型结果 }.recover { case ex => Left(...) // 替换为你的X类型错误结果 } }
关键说明
BroadcastHub.sink(bufferSize):创建一个Sink,将输入流的元素广播给所有订阅它的Source。bufferSize用于设置内部缓冲区大小,可根据流数据量调整。sharedSource是可复用的Source,订阅它的多个消费者会收到完全相同的元素序列。- 流生命周期:当原始输入流
content完成时,sharedSource会向所有下游消费者发送完成信号;单个消费者取消订阅不会影响其他消费者的正常运行。
如果需要更精细的分流控制,也可以使用GraphDSL构建包含Broadcast节点的流图:
private def duplicateStream(content: Source[ByteString, Any]): Future[Either[X, Y]] = { import akka.stream.scaladsl.GraphDSL import akka.stream.ClosedShape val sink = StreamConverters.asInputStream() val graph = RunnableGraph.fromGraph(GraphDSL.create(sink, BroadcastHub.sink[ByteString](256))(Keep.both) { implicit builder => (inputSink, broadcastSink) => import GraphDSL.Implicits._ val broadcast = builder.add(Broadcast[ByteString](2)) content ~> broadcast ~> inputSink broadcast ~> broadcastSink ClosedShape }) val (is, bs) = graph.run() // 后续业务处理逻辑同上 is.map { inputStream => Right(...) }.recover { case ex => Left(...) } }
内容的提问来源于stack exchange,提问作者Sushant Somani
相关产品推荐
相关产品推荐

