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
相关产品推荐
相关产品推荐

