You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Quarkus gRPC大文件上传性能劣化问题求助

大文件gRPC上传性能优化方案(Quarkus Java)

核心问题诊断

  1. 客户端同步阻塞传输:当前代码每次发送一个分片后同步等待响应,完全浪费了gRPC流式传输的异步批量能力,导致大量网络往返延迟。
  2. 事务粒度不合理:服务端流式方法上的@Transactional注解会让每个分片处理都开启/提交事务,高频事务的开销严重拖慢速度。
  3. 冗余数据库操作:每个分片都执行LakeObject的创建/更新,高频数据库IO成为性能瓶颈。
  4. 分片尺寸过小: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 08:04:52