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

gRPC百万消息流式传输内存溢出问题求解

解决gRPC流式传输百万数据时的Netty直接内存溢出问题

你遇到的OutOfDirectMemoryError本质是gRPC异步发送的消息在Netty直接内存中堆积导致的——客户端一直在高速发送消息,而服务器处理速度跟不上,大量未被处理的消息缓冲区占满了Netty的直接内存配额。下面是几个针对性的解决方案:

1. 实现背压(Backpressure)机制(生产环境推荐)

gRPC的StreamObserver原生支持背压逻辑,你可以通过监听onReady()状态控制发送速率,只有当客户端通道确认可以接收更多数据时,再继续发送下一批消息,从根源上避免内存堆积。

修改你的客户端发送逻辑,添加背压控制:

List<MessageValue> millionMessages = new ArrayList<>();
// 优化:复用空MessageValue对象,减少堆内存开销
MessageValue emptyMsg = MessageValue.newBuilder().build();
for (long i = 0; i < 1000000; i++) {
    millionMessages.add(emptyMsg);
}

long before = System.currentTimeMillis();
AtomicInteger sendIndex = new AtomicInteger(0);
final int BATCH_SIZE = 50000;

StreamObserver<MessageValue> requestObserver = new StreamObserver<>() {
    private boolean channelReady = true;

    @Override
    public void onNext(MessageValue value) {}

    @Override
    public void onError(Throwable t) {
        LOG.error("发送失败", t);
    }

    @Override
    public void onCompleted() {
        long total = System.currentTimeMillis() - before;
        LOG.info("百万消息发送完成,总耗时: {}ms", total);
    }

    @Override
    public void onReady() {
        channelReady = true;
        sendNextBatch();
    }

    private void sendNextBatch() {
        while (channelReady && sendIndex.get() < millionMessages.size()) {
            int endIdx = Math.min(sendIndex.get() + BATCH_SIZE, millionMessages.size());
            for (int i = sendIndex.get(); i < endIdx; i++) {
                asyncStub.recievetonm(responseObserverTonMessages).onNext(millionMessages.get(i));
            }
            sendIndex.set(endIdx);
            LOG.info("已发送消息数: {}", endIdx);
            // 暂停发送,等待下一次onReady触发
            channelReady = false;
        }
        // 所有消息发送完毕,标记完成
        if (sendIndex.get() >= millionMessages.size()) {
            asyncStub.recievetonm(responseObserverTonMessages).onCompleted();
        }
    }
};

// 触发首次消息发送
requestObserver.onReady();

2. 临时调整Netty直接内存上限

如果只是临时测试场景,可以通过JVM参数增大Netty的直接内存配额,在启动命令中添加:

-XX:MaxDirectMemorySize=2g

注意:这只是治标方法,无法解决消息堆积的核心问题,不建议在生产环境依赖此方案。

3. 批量打包消息(减少gRPC消息帧数量)

如果你的Protobuf定义允许,可以将多个MessageValue打包成一个批量消息,减少gRPC的消息帧数量,从而降低内存开销:

首先修改Protobuf服务定义:

// 新增批量消息类型
message BatchMessage {
    repeated MessageValue messages = 1;
}

// 修改服务方法为接收批量消息流
service ListMessageSRV {
    rpc recievetonm(stream BatchMessage) returns (YourResponseType);
}

然后客户端批量发送:

List<MessageValue> millionMessages = new ArrayList<>();
MessageValue emptyMsg = MessageValue.newBuilder().build();
for (long i = 0; i < 1000000; i++) {
    millionMessages.add(emptyMsg);
}

long before = System.currentTimeMillis();
StreamObserver<BatchMessage> requestObserver = asyncStub.recievetonm(responseObserverTonMessages);

int BATCH_SIZE = 1000;
for (int i = 0; i < millionMessages.size(); i += BATCH_SIZE) {
    int endIdx = Math.min(i + BATCH_SIZE, millionMessages.size());
    BatchMessage batch = BatchMessage.newBuilder()
            .addAllMessages(millionMessages.subList(i, endIdx))
            .build();
    requestObserver.onNext(batch);
    LOG.info("已发送第{}批,累计消息数: {}", (i/BATCH_SIZE)+1, endIdx);
}

requestObserver.onCompleted();
long total = System.currentTimeMillis() - before;
LOG.info("总耗时: {}ms", total);

以上方案中,背压机制是最适合生产环境的长期解决方案,它能自适应服务器的处理能力,彻底避免内存溢出问题。

内容的提问来源于stack exchange,提问作者Roman Antonio Muñoz Rivas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:25:33