Scala中如何将Source转换为Flow?转换后返回空结果求排查
问题排查与正确实现方案
原代码的核心问题
- 提前调用API导致输入脱节:你在创建Flow之前就执行了
<API CALL>获取output,此时Flow还未接收任何输入数据,API调用根本没用到后续输入的ByteString,返回的InputStream大概率是空的,自然输出为空。 - Flow逻辑完全割裂:
Flow.fromSinkAndSource(sink, mySource)创建的Flow中,输入的ByteString被sink转成InputStream后完全闲置,输出则来自提前创建的mySource,输入和输出没有任何关联,完全不符合你“用输入数据调用API并返回结果”的需求。 - 笔误问题:代码最后一行的
source应为mySource,但即使修正,核心逻辑错误依然存在。
正确实现思路
要实现“输入ByteString流 → 调用API(接收InputStream)→ 返回ByteString流”的Flow,核心是让API调用依赖于输入数据,并将API返回的InputStream转回ByteString流。根据API对输入的要求,分两种场景:
场景1:API需要完整的输入数据(非流式)
如果API必须接收完整的输入内容才能返回结果,先收集所有输入字节,再转成InputStream调用API:
// 假设你的API定义如下 def callApi(inputStream: InputStream): InputStream = { // 处理输入流并返回结果流 } val myFlow: Flow[ByteString, ByteString, NotUsed] = Flow[ByteString] // 收集所有输入字节为单个ByteString .fold(ByteString.empty)(_ ++ _) // 转换为ByteArrayInputStream .map(inputBytes => new ByteArrayInputStream(inputBytes.toArray)) // 调用API获取输出流 .map(callApi) // 将输出InputStream转换为ByteString流 .flatMapConcat(outputStream => StreamConverters.fromInputStream(() => outputStream))
场景2:API支持流式输入输出
如果API可以边接收输入边返回结果,直接将输入ByteString流转为InputStream传给API,再将返回的InputStream转成ByteString流:
def callApi(inputStream: InputStream): InputStream = { // 流式处理输入并返回结果流 } val myFlow: Flow[ByteString, ByteString, NotUsed] = Flow[ByteString] // 将输入ByteString流转换为InputStream,保留Materialized值 .viaMat(StreamConverters.asInputStream())(Keep.right) // 异步调用API(如果API是同步的,用Future.successful包裹;异步则直接返回Future) .mapAsync(1) { inputStream => Future.successful(callApi(inputStream)) } // 将API返回的InputStream转换为ByteString流 .flatMapConcat(outputStream => StreamConverters.fromInputStream(() => outputStream))
关键注意事项
- 确保API调用在Flow的处理逻辑中执行,而非提前初始化,这样才能关联输入数据。
- 如果API是异步的,将
mapAsync中的Future.successful替换为实际的异步API调用返回的Future。 - 使用
flatMapConcat将InputStream转换为Source[ByteString, _],从而嵌入到Flow的输出流中。
内容的提问来源于stack exchange,提问作者Siddharth Shankar
相关产品推荐
相关产品推荐

