分块上传至AWS预签名URL报错:Stream已被监听
问题原因分析
你遇到的Bad state: Stream has already been listened to错误,核心原因是**file.readStream是单订阅流**:这类流只能被监听/消费一次,第一次分块读取后流就会被耗尽并关闭,后续迭代尝试再次操作同一个流时,就会触发该错误。
修复方案
针对这个问题,我们需要每次分块时创建全新的文件流,同时修正代码中的其他潜在问题,具体步骤如下:
1. 替换流的获取方式
利用PlatformFile的path属性(本地文件场景),每次分块时通过File(path)创建新文件实例,调用openRead(start, end)直接获取对应分块的字节流,确保每次都是全新的可监听流。
如果是Web平台,path属性会为空,需要先将文件内容缓存到内存,再从缓存中截取分块生成流。
2. 修正分块大小与编号
- 创建
MultipartFile时,必须传入当前分块的实际大小(end - start),而非整个文件大小,否则AWS服务端会校验失败。 - AWS S3分块上传的
PartNumber必须从1开始递增,循环中需用i+1作为当前分块的编号,而非固定参数值。
3. 清理无效代码
删除_calculateStreamSize方法:该方法通过异步监听流计算大小,不仅返回值始终为0,还会提前耗尽流,完全无实际作用。
修改后的关键代码
import 'dart:io'; import 'dart:typed_data'; import 'package:flutter/foundation.dart'; import 'package:http/http.dart' as http; import 'package:file_picker/file_picker.dart'; class UploadRequest { String url; final String method; final String fileKey; final Map<String, String>? bodyData; final Map<String, String>? headers; final PlatformFile file; final Function(int)? onUploadProgress; late final int _maxChunkSize; int fileSize; String fileName; Map<String, String>? uploadHeader; List<Map<String, dynamic>>? tags; late final Uint8List? _fileBytes; // Web平台文件缓存 UploadRequest({ required this.url, this.method = "POST", this.fileKey = "file", this.bodyData = const {}, required this.file, this.onUploadProgress, required int maxChunkSize, required this.fileSize, required this.fileName, this.headers, this.uploadHeader, this.tags, }) { _maxChunkSize = min(fileSize, maxChunkSize); // Web平台预读取文件到内存 if (kIsWeb) { _fileBytes = file.bytes ?? file.readStream!.expand((chunk) => chunk).toListSync(); } else { _fileBytes = null; } } Future<void> upload() async { List<Map<String, dynamic>> _tag = []; for (int i = 0; i < _chunksCount; i++) { print("Uploading chunk $i"); final start = _getChunkStart(i); final end = _getChunkEnd(i); final chunkSize = end - start; var chunkStream = _getChunkStream(start, end); var request = http.MultipartRequest( method, Uri.parse(url), ); // 添加请求头,包含Content-Range request.headers.addAll(_getHeaders(start, end)); request.fields.addAll(bodyData!); request.files.add(http.MultipartFile( fileKey, chunkStream, chunkSize, // 传入当前分块实际大小 filename: fileName, )); var response = await request.send(); String result = await response.stream.bytesToString(); uploadHeader = response.headers; // 使用i+1作为分块编号(AWS要求从1开始) _tag.add({"etag": uploadHeader!["etag"], "PartNumber": i + 1}); tags = _tag; print("Chunk $i upload result: $result"); print("Response headers: ${response.headers}"); _updateProgress(i, chunkSize, fileSize); } } Stream<List<int>> _getChunkStream(int start, int end) { if (kIsWeb) { if (_fileBytes == null) { throw StateError("Web platform: file bytes not available"); } // 从内存缓存中截取分块并转为流 final chunkBytes = _fileBytes!.sublist(start, end); return Stream.value(chunkBytes); } else { if (file.path == null) { throw StateError("Local file path is null"); } // 本地文件直接创建指定范围的流 return File(file.path!).openRead(start, end); } } // 上传进度更新 void _updateProgress(int chunkIndex, int chunkCurrent, int chunkTotal) { int totalUploadedSize = (chunkIndex * _maxChunkSize) + chunkCurrent; double totalUploadProgress = (totalUploadedSize / fileSize) * 100; onUploadProgress?.call(totalUploadProgress.truncate()); } // 分块起始位置 int _getChunkStart(int chunkIndex) => chunkIndex * _maxChunkSize; // 分块结束位置 int _getChunkEnd(int chunkIndex) => min((chunkIndex + 1) * _maxChunkSize, fileSize); // 生成包含Content-Range的请求头 Map<String, String> _getHeaders(int start, int end) { var header = {'Content-Range': 'bytes $start-${end - 1}/$fileSize'}; if (headers != null) { header.addAll(headers!); } return header; } // 计算分块总数 int get _chunksCount { var result = (fileSize / _maxChunkSize).ceil(); return result; } }
内容的提问来源于stack exchange,提问作者Naveen
相关产品推荐
相关产品推荐

