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

GRPC在Spigot插件中频繁创建线程致服务器崩溃问题排查

问题根源分析与解决方案

这个问题完全是由你的代码逻辑导致的,和gRPC本身无关——错误的线程调度方式导致大量阻塞线程持续堆积,最终耗尽了进程的文件句柄上限。

为什么会出现这个问题?

我们来拆解一下你的代码逻辑和gRPC阻塞式Stub的行为:

  1. gRPC阻塞式Stub的特性:你使用的ChatSyncBlockingStub返回的Iterator<ChatMessage>是阻塞式的——调用hasNext()或next()时,如果没有新消息到来,当前线程会一直进入等待状态(也就是你线程堆栈里看到的WAITING状态),直到有新消息或者连接断开。
  2. Bukkit定时任务的行为:你的GRPCMessageReceiver继承了BukkitRunnable,如果是通过runTaskTimer这类方法调度的,Bukkit会每隔指定周期就从线程池里启动一个新线程来执行run()方法。
  3. 线程堆积的过程:第一次执行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 17:42:36