Alpakka 4.0.0中S3.getObject截断文件内容的问题求助
问题根源与解决方案
你的代码出现文件内容截断的核心原因是对Alpakka S3的getObject返回值和流式操作的理解偏差:
原代码的问题点
S3.getObject直接返回的是整个文件内容的流式数据源(Source[ByteString, NotUsed]),而非单个结果;- 你使用
.zip(Source.single(r.key))时,zip操作会严格按两个流的元素一一配对,当Source.single(r.key)发出唯一元素后,流会立即终止,导致文件内容的流只取到第一个ByteString块就被截断。
正确实现方式
需要将文件流中的每一块字节数据都与对应的文件名绑定,而非将整个文件流和单个文件名做配对。以下是修正后的代码:
方案1(保留并行下载配置)
def downloadFilesSource(bucketName: String, primaryKey: String)( implicit ec: ExecutionContext ): Source[(ByteString, String), NotUsed] = { S3.listBucket(bucketName, prefix = Some(primaryKey)) // 并行获取每个文件的流(保持原有的并行度配置) .mapAsync(configuration.config.downloadParallelism.value) { r => // getObject是惰性创建流,无需额外Future包裹,直接返回即可 Future.successful( S3.getObject(r.bucketName, r.key).map(byteStr => (byteStr, r.key)) ) } // 将每个文件的流合并为一个连续的输出流 .flatMapConcat(identity) }
方案2(更简洁的无并行写法)
如果不需要并行初始化文件流,可直接用map替代mapAsync:
def downloadFilesSource(bucketName: String, primaryKey: String): Source[(ByteString, String), NotUsed] = { S3.listBucket(bucketName, prefix = Some(primaryKey)) .map { r => // 给文件流的每一块数据都带上对应的文件名 S3.getObject(r.bucketName, r.key).map((_, r.key)) } .flatMapConcat(identity) }
为什么这样能解决问题
S3.getObject(...).map((_, r.key))会遍历文件流的所有ByteString块,为每个字节块绑定当前文件的key,确保整个文件的所有内容都会被输出;flatMapConcat会将多个文件的流合并成一个连续的流,既保留了原有的流式特性,又不会因配对逻辑截断内容。
内容的提问来源于stack exchange,提问作者user3468054
相关产品推荐
相关产品推荐

