Quarkus gRPC大文件上传性能劣化问题求助
大文件gRPC上传性能优化方案(Quarkus Java)
核心问题诊断
- 客户端同步阻塞传输:当前代码每次发送一个分片后同步等待响应,完全浪费了gRPC流式传输的异步批量能力,导致大量网络往返延迟。
- 事务粒度不合理:服务端流式方法上的
@Transactional注解会让每个分片处理都开启/提交事务,高频事务的开销严重拖慢速度。 - 冗余数据库操作:每个分片都执行LakeObject的创建/更新,高频数据库IO成为性能瓶颈。
- 分片尺寸过小:64KB分片导致20MB文件需要320次请求,大幅增加gRPC协议的 overhead。
具体优化方案
1. 客户端改为异步流式批量发送
将双向流从“多次同步单请求”改为真正的流式传输,仅第一个分片需要等待cid,后续分片异步批量发送,无需逐个等待响应:
ObjectModel obj = new ObjectModel(); obj.cid = "null"; try { // 处理第一个分片,获取初始CID byte[] firstBytes = new byte[1024 * 1024]; // 先调整为1MB分片 int firstLen = input.is.read(firstBytes); if (firstLen == -1) { // 空文件处理逻辑 return; } ByteString firstByteString = ByteString.copyFrom(firstBytes, 0, firstLen); String cid = coreChunkGrpcClient.uploadCoreChunk(Multi.createFrom().item(UploadCoreChunkRequest.newBuilder() .setBearer("Bearer " + jwt.getRawToken()) .setOutput(io.quarkus.example.ObjectFormModelChunk.newBuilder() .setMetadata(output.metadata) .setIs(firstByteString) .setFileID(obj.cid) .build()) .build())) .onItem().transform(UploadCoreResponse::getToken) .await().indefinitely(); obj.cid = cid; // 异步批量发送后续分片,无需逐个等待响应 byte[] bytes = new byte[1024 * 1024]; int len; Multi<UploadCoreChunkRequest> chunkMulti = Multi.createFrom().generator(() -> input.is, (is, emitter) -> { len = is.read(bytes); if (len == -1) { emitter.complete(); return; } ByteString byteString = ByteString.copyFrom(bytes, 0, len); emitter.emit(UploadCoreChunkRequest.newBuilder() .setBearer("Bearer " + jwt.getRawToken()) .setOutput(io.quarkus.example.ObjectFormModelChunk.newBuilder() .setMetadata(output.metadata) .setIs(byteString) .setFileID(cid) .build()) .build()); }); // 异步订阅,仅处理错误和完成事件 coreChunkGrpcClient.uploadCoreChunk(chunkMulti) .subscribe().with( response -> {}, // 无需处理每个分片响应 error -> log.error("Chunk upload failed", error), () -> log.info("All chunks uploaded successfully") ); } catch (IOException e) { log.error("File read error", e); }
2. 优化服务端事务与数据库操作
- 移除流式方法上的
@Transactional,仅在必要的数据库操作处添加事务 - 仅第一个分片执行数据库插入,后续分片跳过冗余操作:
// 服务端uploadCoreChunk方法:移除@Transactional注解 @Override public StreamObserver<UploadCoreChunkRequest> uploadCoreChunk(StreamObserver<UploadCoreResponse> responseObserver) { return new StreamObserver<UploadCoreChunkRequest>() { private String currentCid; @Override public void onNext(UploadCoreChunkRequest request) { var bearer = request.getBearer(); var inputs = request.getOutput(); String cid = inputs.getFileID(); try { var result = insertFileChunk(bearer, inputs, cid); currentCid = result.getResp().toString(); // 可选:减少响应次数,比如每10个分片返回一次 responseObserver.onNext(UploadCoreResponse.newBuilder() .setCode(result.getCode()) .setMessage(result.getMsg()) .setToken(currentCid) .build()); } catch (IOException e) { log.error("Error processing chunk", e); responseObserver.onNext(UploadCoreResponse.newBuilder() .setCode(500) .setMessage("Internal Error") .build()); } } @Override public void onError(Throwable t) { log.error("Upload stream error", t); responseObserver.onError(t); } @Override public void onCompleted() { responseObserver.onNext(UploadCoreResponse.newBuilder() .setCode(200) .setMessage("All chunks uploaded") .setToken(currentCid) .build()); responseObserver.onCompleted(); } }; } // insertFileChunk方法:仅第一个分片执行数据库操作 @Transactional // 仅在需要事务的地方添加 public LakeHttpResponse insertFileChunk(String bearer, ObjectFormModelChunk input, String cid) throws IOException { ByteString is = input.getIs(); if (cid == null || cid.isEmpty() || cid.isBlank() || cid.equals("null")) { LakeObjectMetadata meta = parseMetadata(input.getMetadata()); // 补充元数据解析逻辑 cid = fs.create(meta.getName(), meta.getLength(), is); // 仅第一个分片创建数据库记录 LakeObject object = new LakeObject(); object.setCid(cid); Long now = new Date().getTime(); object.setCreateTime(now); object.setAccessTime(now); object.setParentId(0L); entityManager.persist(object); log.info("Object metadata saved, cid={}", cid); return new LakeHttpResponse(200, "OK", cid); } else { if (!fs.insertChunk(cid, is)) { return new LakeHttpResponse(500, "Internal error", null); } // 后续分片跳过数据库操作 return new LakeHttpResponse(200, "OK", cid); } }
3. 调整分片尺寸与内存拷贝优化
- 将分片大小从64KB调整为1MB(可根据DFS块大小进一步调整,建议不超过4MB),减少请求次数
- 直接从InputStream生成ByteString,减少中间内存拷贝:
// 替换原ByteSting.copyFrom方式,直接读取输入流 ByteString byteString = ByteString.readFrom(new ByteArrayInputStream(bytes, 0, len));
4. gRPC配置优化
在application.properties中添加以下配置,提升连接和流性能:
# 客户端配置 quarkus.grpc.clients.coreChunkGrpcClient.max-inbound-message-size=67108864 # 64MB quarkus.grpc.clients.coreChunkGrpcClient.max-retry-attempts=0 # 关闭大文件重试 quarkus.grpc.clients.coreChunkGrpcClient.keep-alive-time=30s quarkus.grpc.clients.coreChunkGrpcClient.keep-alive-timeout=5s # 服务端配置 quarkus.grpc.server.max-inbound-message-size=67108864 quarkus.grpc.server.max-concurrent-streams=100 # 增大并发流数
内容的提问来源于stack exchange,提问作者midmash36
相关产品推荐
相关产品推荐

