GRPC在Spigot插件中频繁创建线程致服务器崩溃问题排查
问题根源分析与解决方案
这个问题完全是由你的代码逻辑导致的,和gRPC本身无关——错误的线程调度方式导致大量阻塞线程持续堆积,最终耗尽了进程的文件句柄上限。
为什么会出现这个问题?
我们来拆解一下你的代码逻辑和gRPC阻塞式Stub的行为:
- gRPC阻塞式Stub的特性:你使用的
ChatSyncBlockingStub返回的Iterator<ChatMessage>是阻塞式的——调用hasNext()或next()时,如果没有新消息到来,当前线程会一直进入等待状态(也就是你线程堆栈里看到的WAITING状态),直到有新消息或者连接断开。 - Bukkit定时任务的行为:你的
GRPCMessageReceiver继承了BukkitRunnable,如果是通过runTaskTimer这类方法调度的,Bukkit会每隔指定周期就从线程池里启动一个新线程来执行run()方法。 - 线程堆积的过程:第一次执行
run()时,你初始化了receivingIterator并调用hasNext(),线程开始阻塞等待消息;下一个调度周期到了,Bukkit又启动新线程执行run(),此时receivingIterator已经不为null,新线程再次调用hasNext()并阻塞。就这样,每过一个调度周期就新增一个阻塞线程,线程数量指数级增长,最终耗尽了系统分配给进程的文件句柄(每个线程都会占用系统资源,包括文件句柄)。
怎么解决这个问题?
核心思路是:让单个线程持续监听gRPC流,而不是通过定时任务重复启动新线程。这里提供两种可行的方案:
方案1:使用单独线程处理阻塞式gRPC流
放弃Bukkit的定时调度,自己创建一个独立线程来持续监听gRPC流,同时注意Bukkit API的线程安全问题(Bukkit API只能在主线程调用):
public class GRPCMessageReceiver implements Runnable { private final ChatSyncGrpc.ChatSyncBlockingStub chatSync; private final LockedQueue lockedQueue; private final Plugin yourPlugin; // 传入你的插件实例 private volatile boolean isRunning = true; public GRPCMessageReceiver(ChatSyncGrpc.ChatSyncBlockingStub chatSync, LockedQueue queue, Plugin plugin) { this.chatSync = chatSync; this.lockedQueue = queue; this.yourPlugin = plugin; // 启动独立监听线程 new Thread(this, "GRPC-ChatSync-Listener").start(); } @Override public void run() { while (isRunning) { try { // 建立gRPC流连接 Iterator<ChatMessage> messageIterator = chatSync.beginReceive(Empty.newBuilder().build()); // 持续读取消息 while (messageIterator.hasNext()) { ChatMessage message = messageIterator.next(); // 将消息处理逻辑提交到Bukkit主线程执行 Bukkit.getScheduler().runTask(yourPlugin, () -> { lockedQueue.Put(message); // 这里添加你的消息处理逻辑,比如广播到游戏内 }); } } catch (StatusRuntimeException e) { Bukkit.getLogger().info("Chat Sync连接断开,5秒后尝试重连..."); // 重连前休眠,避免无限循环重试 try { Thread.sleep(5000); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } } } // 插件卸载时调用,停止监听线程 public void stopListening() { isRunning = false; } }
方案2:改用gRPC异步Stub(更推荐)
gRPC的异步Stub(ChatSyncStub)配合StreamObserver观察者模式,不需要阻塞线程,更适合Bukkit的异步环境:
public class GRPCMessageReceiver { private final ChatSyncGrpc.ChatSyncStub chatSync; private final LockedQueue lockedQueue; private final Plugin yourPlugin; private StreamObserver<Empty> requestObserver; public GRPCMessageReceiver(ChatSyncGrpc.ChatSyncStub chatSync, LockedQueue queue, Plugin plugin) { this.chatSync = chatSync; this.lockedQueue = queue; this.yourPlugin = plugin; // 启动异步监听 startAsyncListening(); } private void startAsyncListening() { StreamObserver<ChatMessage> responseObserver = new StreamObserver<>() { @Override public void onNext(ChatMessage message) { // 提交到Bukkit主线程处理消息 Bukkit.getScheduler().runTask(yourPlugin, () -> { lockedQueue.Put(message); // 处理消息的逻辑 }); } @Override public void onError(Throwable t) { Bukkit.getLogger().info("Chat Sync连接出错,5秒后重连..."); try { Thread.sleep(5000); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); return; } // 重新启动监听 startAsyncListening(); } @Override public void onCompleted() { Bukkit.getLogger().info("Chat Sync流已结束,5秒后重连..."); startAsyncListening(); } }; // 发起异步流请求 requestObserver = chatSync.beginReceive(responseObserver); requestObserver.onNext(Empty.newBuilder().build()); requestObserver.onCompleted(); } // 插件卸载时调用,关闭流 public void stopListening() { if (requestObserver != null) { requestObserver.onCompleted(); } } }
额外注意事项
- 确保你的
LockedQueue是线程安全的,或者在主线程中处理读写(因为上面的方案都是把消息提交到主线程写入队列)。 - 在你的插件的
onDisable()方法中,一定要调用GRPCMessageReceiver的stopListening()方法,避免资源泄漏。 - 根据你的业务场景调整重连的休眠时间,避免频繁重试浪费资源。
内容的提问来源于stack exchange,提问作者ForberichN
相关产品推荐
相关产品推荐

