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

Android gRPC聊天应用发送首条消息报cancelled before half close异常

gRPC双向流聊天应用CANCELLED: cancelled before receiving half close异常修复

这个异常的核心触发逻辑是:双向RPC流的某一端没有完成正常的半关闭(half close)握手流程就意外断开,而另一端还持有这个失效流的引用,尝试向已经被取消的流写入数据,就会抛出该错误。

问题根因

  • 客户端线程违规:gRPC异步Stub的回调(onNext/onError/onCompleted)默认运行在gRPC内部的网络线程,你直接在回调里调用Adapter更新UI、刷新RecyclerView,违反Android主线程唯一操作UI的规则,会触发CalledFromWrongThreadException导致客户端流意外崩溃断开。断开时没有给服务端发送正常的流结束帧,服务端无法第一时间感知连接失效。这也是为什么第一条消息能正常广播,之后立刻抛异常——客户端收到第一条消息的回调里操作UI崩溃,流被强制终止,服务端感知到流被取消时还持有引用,就抛出对应错误。
  • 服务端连接管理缺陷:你存储客户端连接的observers集合是普通非线程安全的ArrayList,多客户端并发连接、发消息、断开时,遍历、新增、删除操作会触发并发修改问题,也无法及时剔除已经失效的连接引用。
  • 缺失异常兜底:服务端广播消息时没有对单个连接的写入操作做异常捕获,单个死连接的报错会直接中断整个广播流程,甚至影响其他正常连接。
  • 客户端缺失生命周期回收:Activity销毁(旋转屏幕、返回退出)时没有正常关闭流和Channel,会产生大量僵尸连接留在服务端的observer集合里。

修复方案

服务端代码修改

将observers替换为线程安全集合,广播时增加异常捕获自动剔除失效连接,补全流关闭的资源清理逻辑:

// 首先把observers全局定义替换为线程安全集合
private final List<StreamObserver<Messaging.Message>> observers = new CopyOnWriteArrayList<>();

@Override
public StreamObserver<Messaging.Message> sendReceiveMessage(StreamObserver<Messaging.Message> responseObserver) {
    observers.add(responseObserver);

    return new StreamObserver<Messaging.Message>() {
        @Override
        public void onNext(Messaging.Message message) {
            System.out.println(String.format("Got a message from: '%s' : '%s'", message.getMessageOwnerId(), message.getMessage()));
            // 广播时捕获单连接异常,自动剔除失效连接
            for (StreamObserver<Messaging.Message> o : observers) {
                try {
                    o.onNext(message);
                } catch (Exception e) {
                    observers.remove(o);
                    o.onError(Status.CANCELLED.withCause(e).asRuntimeException());
                }
            }
        }

        @Override
        public void onError(Throwable throwable) {
            observers.remove(responseObserver);
            responseObserver.onError(throwable);
        }

        @Override
        public void onCompleted() {
            observers.remove(responseObserver);
            responseObserver.onCompleted();
        }
    };
}

客户端代码修改

所有UI操作切回主线程执行,补全Activity销毁时的资源释放逻辑:

// 新增主线程Handler用于线程切换
private final Handler mainHandler = new Handler(Looper.getMainLooper());

private void createStub(){
    channel = ManagedChannelBuilder.forAddress("10.0.2.2",8080)
            .usePlaintext()
            .build();

    stub = MessagingServiceGrpc.newStub(channel);

    toServer = stub.sendReceiveMessage(new StreamObserver<Messaging.Message>() {
        @Override
        public void onNext(Messaging.Message value) {
            // 切到主线程再操作UI组件
            mainHandler.post(() -> {
                msgAdapter.receiveMessage(value);
                Log.d("ChatMsg",value.getMessageOwnerId() + " : " + value.getMessage() );
                msgAdapter.notifyDataSetChanged();
            });
        }

        @Override
        public void onError(Throwable t) {
            t.printStackTrace();
            // 可按需添加流断开后的延迟重连逻辑,兼容网络波动场景
        }

        @Override
        public void onCompleted() {
            Log.d("ChatMsg","Stream closed normally");
        }
    });
}

// 重写onDestroy方法释放gRPC资源
@Override
protected void onDestroy() {
    super.onDestroy();
    // 先正常发送半关闭信号
    if (toServer != null) {
        try {
            toServer.onCompleted();
        } catch (Exception ignored) {}
    }
    // 关闭通信通道
    if (channel != null && !channel.isShutdown()) {
        channel.shutdown();
    }
    // 移除主线程所有待执行回调,避免内存泄漏
    mainHandler.removeCallbacksAndMessages(null);
}

额外优化建议

  • 可以给每个连接绑定客户端唯一ID,避免同一个客户端重复创建连接产生冗余observer;
  • 添加轻量心跳机制,定时检测连接活性,提前剔除僵尸连接;
  • 消息广播可以按chat_id做分组,不要给所有连接全量广播,减少无效请求。

内容的提问来源于stack exchange,提问作者Davidflow

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 01:09:34