一、从一次线上事故说起

有一天下午,运营同事跑来喊:“列表页打不开了!”我第一反应是看监控,结果发现订单服务、用户服务、库存服务的错误率全在往上窜。奇怪的是,这些服务的CPU和内存都没满,但请求就是卡住,最后全部超时。再看日志,里面铺满了 ThreadPoolExecutor 的拒绝异常,还有一堆 TimeoutException。顺着调用链往上查,发现所有链路都汇聚到一个gRPC服务上,这个服务的线程池大小是200,可活跃线程数一直是200,队列也满了。说白了,线程池被占满了,新请求进不来,老请求也在里面干等,整个服务就僵在那儿了。那之后我们连续排查了好多天,发现问题的根源并不只是简单的流量大,而是一连串小毛病叠加在一起,最终引发了这台“机器”的全面瘫痪。今天咱们就从这个事故出发,把gRPC服务端线程池耗尽的来龙去脉彻底捋清楚,再聊聊怎么从根上避免这种问题。

二、先弄懂gRPC服务端线程模型

2.1 线程池是怎么来的

gRPC服务端在启动的时候,会创建一组线程用来处理网络请求。咱们常用的Java技术栈里,grpc-java底层用的是Netty,Netty的EventLoop线程负责读写网络数据,而业务逻辑默认是在专门的线程池里执行的。这个线程池就是“服务端线程池”,它的大小可以配置,也可以使用默认值。默认情况下,grpc-java会创建一个固定大小的线程池,数量等于CPU核数乘以一个系数。生产环境里大家往往手动调大,比如设置成200、500甚至更多,以为这样就能扛住更大的并发。

这个线程池就像一个“营业厅的窗口”。每个请求进来,就有一个窗口工作人员去接待。如果所有窗口都在忙,后面的客户就只能排队,排队的区域就是阻塞队列。再满的话,新来的客户直接被告知“今天不营业了”,对应线程池的拒绝策略。

2.2 一个请求的完整旅程

为了直观理解,咱们看一个最小的gRPC服务端代码,技术栈是Java,使用grpc-java和protobuf。先定义一个简单的接口。

// 文件: hello.proto
syntax = "proto3";

package demo;

service Greeter {
  // 一个简单的打招呼接口
  rpc SayHello (HelloRequest) returns (HelloReply) {}
}

message HelloRequest {
  string name = 1;
}

message HelloReply {
  string message = 1;
}

然后实现服务端逻辑。这里故意让处理逻辑里睡2秒,模拟一个需要长时间等待的操作。

// 文件: GreeterServiceImpl.java
package demo;

import io.grpc.stub.StreamObserver;

// 服务实现类,每个请求都会占用一个线程池中的线程
public class GreeterServiceImpl extends GreeterGrpc.GreeterImplBase {

    @Override
    public void sayHello(HelloRequest request, StreamObserver<HelloReply> responseObserver) {
        // 线程池线程正在处理这个请求
        System.out.println("收到请求: " + request.getName() + ",当前线程: "
                + Thread.currentThread().getName());

        try {
            // 模拟一个耗时的阻塞操作,比如查数据库、调外部接口
            Thread.sleep(2000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }

        // 构造响应并返回
        HelloReply reply = HelloReply.newBuilder()
                .setMessage("你好, " + request.getName())
                .build();
        responseObserver.onNext(reply);
        responseObserver.onCompleted();
    }
}

再看服务端启动类和线程池配置。grpc-java里可以用 ServerBuilder 直接配置执行器,也就是那个“服务端线程池”。这里我们故意把线程池设成只有2个线程,方便做演示。

// 文件: GrpcServerMain.java
package demo;

import io.grpc.Server;
import io.grpc.ServerBuilder;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class GrpcServerMain {

    public static void main(String[] args) throws Exception {
        // 创建一个只有2个线程的线程池,模拟“窗口很少”的场景
        ExecutorService executor = Executors.newFixedThreadPool(2);

        // 构建gRPC服务,监听9000端口,并指定用这个2线程池
        Server server = ServerBuilder.forPort(9000)
                .addService(new GreeterServiceImpl())
                .executor(executor)
                .build()
                .start();

        System.out.println("gRPC服务已启动,端口: 9000");
        // 进程不退出,等待请求
        server.awaitTermination();
    }
}

这段代码里,只要同一时刻有超过2个请求进来,多出来的请求就会在队列里排队。如果处理一个请求要2秒,那第3个请求就得等前面某个线程释放。看起来“排队”还能接受,但排队时间一长,客户端就会超时,然后在客户端那里引发新的重试,进一步加重服务端负担。

三、阻塞调用是如何把线程池堵死的

3.1 同步阻塞是最常见的坑

很多业务代码都是“一条道走到黑”的写法:gRPC接口收到请求以后,同步去调用数据库、Redis、别的HTTP服务,等拿到结果后才返回。比如下面这段代码,就调用了另一个阻塞式的HTTP客户端。

// 文件: GreeterServiceImpl.java(改造版,引入外部调用)
package demo;

import io.grpc.stub.StreamObserver;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;

public class GreeterServiceImpl extends GreeterGrpc.GreeterImplBase {

    // 一个普通的Java HttpClient,默认是同步阻塞的
    private final HttpClient httpClient = HttpClient.newHttpClient();

    @Override
    public void sayHello(HelloRequest request, StreamObserver<HelloReply> responseObserver) {
        System.out.println("开始处理请求,线程: " + Thread.currentThread().getName());

        try {
            // 模拟调用下游服务,这个调用会阻塞当前线程
            String result = callDownstream(request.getName());

            HelloReply reply = HelloReply.newBuilder()
                    .setMessage("下游返回: " + result)
                    .build();
            responseObserver.onNext(reply);
            responseObserver.onCompleted();
        } catch (Exception e) {
            responseObserver.onError(e);
        }
    }

    // 一个模仿真实阻塞调用的方法
    private String callDownstream(String name) throws Exception {
        HttpRequest httpRequest = HttpRequest.newBuilder()
                .uri(URI.create("http://localhost:8080/api/hello?name=" + name))
                .GET()
                .build();
        // 这行代码会一直等到下游返回,期间线程什么都不干,就等着
        HttpResponse<String> response = httpClient.send(httpRequest, HttpResponse.BodyHandlers.ofString());
        return response.body();
    }
}

注意,httpClient.send 是一个同步方法,它在等待网络响应的过程中,底层线程会进入阻塞状态。这个线程不能干别的活,就干等着。假如下游服务变慢,比如响应时间从20毫秒变成3秒,那这个gRPC服务线程就被“粘”住了。原来一个线程每秒能处理几十个请求,现在几秒才能处理一个。如果线程池一共就200个线程,那只要同时有200个请求卡在下游,整个线程池就全部被占满,后续请求全部排队甚至被拒绝。

3.2 线程耗尽后的连锁反应

线程池一旦耗尽,新来的请求会进入阻塞队列。如果队列也满了,就会触发拒绝策略。默认的 AbortPolicy 会直接抛出 RejectedExecutionException,grpc-java会把这个异常包装成一个 StatusRuntimeException 返回给客户端。客户端收到这种异常后,很多框架会自动重试。于是,客户端疯狂重试,服务端线程池依旧没位置,重试请求也进不来,又加剧了排队。同时,那些在队列里等待的请求,又因为等太久,客户端等不及就超时断开了连接。服务端线程处理一个已经断开的请求,白白浪费资源。

更可怕的是,故障会顺着调用链传播。比如服务A调服务B,服务B线程池满了,B的请求在排队,A的线程也在等待B的响应,于是A的线程池也被慢慢占满。这就是“线程池耗尽引发的连锁故障”。用生活的话说,就是一家餐厅的厨房出菜慢了,所有的服务员都在等菜,没工夫招呼新客人,然后新客人也饿着,整个餐厅乱成一团。

四、操作系统调度异常:线程过多引发的系统性问题

当线程池快满的时候,有些团队会想“那就把线程池调大啊”。于是把线程数从200调到1000,甚至更多。但这往往会让事情变得更糟,因为操作系统调度上千个线程也是一笔巨大的开销。咱们简单说一下操作系统调度是什么。每个线程都是一个执行流,CPU核心数有限,比如8核,那同时能运行的线程就只有8个。其他线程都在等待被调度。操作系统需要频繁地切换线程,把当前运行的线程保存起来,再把下一个线程的数据加载进来,这个过程叫做“上下文切换”。上下文切换是有代价的,它消耗CPU时间,而且切换得太频繁,CPU会一直在做“换人”的工作,而不是真正干业务。

如果服务端线程池里有2000个线程,但大部分都在等待网络IO,那CPU就会大量花费在管理这些线程上。再加上那些被阻塞的线程持有锁、数据库连接、内存资源,整个进程的稳定性都会下降。比如监控指标里会出现线程数量暴涨、CPU使用率忽高忽低、GC频率上升、甚至出现“活锁”一样的现象。咱们来看一个极端的例子:假如我们把线程池配置成直接使用无界队列,并且最大线程数设成几千,那么当流量高峰到来时,线程会不断创建,内存飙升,最终导致 OutOfMemoryError。这种情况下,操作系统可能需要杀掉进程,整个服务彻底宕机。

所以,线程池不是越大越好。这里我们要明白一个核心原理:如果线程要做的是IO等待,那很多阻塞线程其实是“假负载”。与其让它们傻等,不如让它们去做别的事,或者干脆把调用变成异步的。这就像餐厅服务员不应该站在后厨等菜,而是先去招呼其他客人,等菜好了再端过去。

五、完整的全局治理策略

5.1 调用侧改造:异步非阻塞

从根源上解决线程池耗尽,就要尽量避免“阻塞调用”。在Java生态里,我们可以用异步方式发起调用,比如使用 CompletableFuture 或者专门的异步HTTP客户端。gRPC本身也支持异步Stub,咱们可以把同步业务改成异步组合。下面是一个简单的改造示例。

// 文件: AsyncGreeterServiceImpl.java
package demo;

import io.grpc.stub.StreamObserver;
import java.util.concurrent.CompletableFuture;

public class AsyncGreeterServiceImpl extends GreeterGrpc.GreeterImplBase {

    @Override
    public void sayHello(HelloRequest request, StreamObserver<HelloReply> responseObserver) {
        // 不阻塞当前线程,异步发起“查询操作”
        CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
            // 这里模拟一个耗时的外部调用,比如数据库查询
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            return "你好, " + request.getName();
        });

        // 当异步任务完成时,再通过StreamObserver返回结果
        future.whenComplete((result, error) -> {
            if (error != null) {
                responseObserver.onError(error);
                return;
            }
            HelloReply reply = HelloReply.newBuilder()
                    .setMessage(result)
                    .build();
            responseObserver.onNext(reply);
            responseObserver.onCompleted();
        });
    }
}

注意,supplyAsync 默认用的公共ForkJoinPool,这个池子也可能被阻塞任务堵住。更稳妥的做法是给异步操作单独分配一个专门的线程池,并调整大小。但至少业务线程不会被一个请求“粘住”好几秒。如果你用Spring WebFlux或者Vert.x,就能实现真正的全链路异步,但这里咱们不展开。

5.2 服务侧治理:合理的线程池配置

既然不能无限加线程,那就得给线程池“立规矩”。我们可以自己创建线程池,设置核心线程数、最大线程数、队列容量、拒绝策略。下面是一个更合理的配置示例。

// 文件: ThreadPoolConfig.java
package demo;

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class ThreadPoolConfig {

    // 创建一个适合IO密集型的线程池,有界队列,自定义拒绝策略
    public static ThreadPoolExecutor createGrpcExecutor() {
        int coreSize = 50;      // 核心线程数
        int maxSize = 100;      // 最大线程数
        int queueCapacity = 200; // 有界队列容量

        // 线程工厂,给线程起一个有意义的名字,方便排查问题
        ThreadFactory threadFactory = new ThreadFactory() {
            private final AtomicInteger count = new AtomicInteger(1);

            @Override
            public Thread newThread(Runnable r) {
                Thread t = new Thread(r, "grpc-worker-" + count.getAndIncrement());
                t.setDaemon(true);
                return t;
            }
        };

        return new ThreadPoolExecutor(
                coreSize,
                maxSize,
                30, // 空闲线程存活时间
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(queueCapacity),
                threadFactory,
                // 自定义拒绝策略:调用线程执行,或者打个日志,而不是简单抛异常
                new ThreadPoolExecutor.CallerRunsPolicy()
        );
    }
}

这里用 CallerRunsPolicy 作为拒绝策略。它的意思是,如果线程池和队列都满了,新任务不会被丢弃,而是由提交任务的线程自己去执行。这有一个优点:不会丢失请求,同时起到了天然的限流作用(比如Netty的EventLoop线程被占用后,网络读取也会变慢)。但也有一个缺点:如果任务本身很耗时,可能会阻塞Netty的EventLoop线程,影响网络读取。所以这个策略要看场景,也可以自定义一个策略,比如把任务放入一个持久化队列,或者直接返回一个“稍后再试”的响应。

5.3 流量控制与隔离

线程池治理不能只看一个服务,还需要在整体架构上做流量控制。流量控制就像给每个入口加一个“限流闸门”。我们可以用令牌桶算法,或者简单的信号量,限制最大并发数。比如用Java的 Semaphore 控制同时处理的请求数。

// 文件: RateLimitedService.java
package demo;

import io.grpc.stub.StreamObserver;
import java.util.concurrent.Semaphore;

public class RateLimitedService extends GreeterGrpc.GreeterImplBase {

    // 最大同时处理5个请求,多出来的直接拒绝
    private final Semaphore semaphore = new Semaphore(5);

    @Override
    public void sayHello(HelloRequest request, StreamObserver<HelloReply> responseObserver) {
        // 尝试获取许可,获取不到就快速返回失败
        if (!semaphore.tryAcquire()) {
            responseObserver.onError(
                    io.grpc.Status.RESOURCE_EXHAUSTED
                            .withDescription("系统繁忙,请稍后重试")
                            .asRuntimeException()
            );
            return;
        }

        try {
            // 模拟业务处理
            String result = "接收请求: " + request.getName();
            HelloReply reply = HelloReply.newBuilder()
                    .setMessage(result)
                    .build();
            responseObserver.onNext(reply);
            responseObserver.onCompleted();
        } finally {
            // 处理完成,释放许可
            semaphore.release();
        }
    }
}

服务隔离也很重要。不要把消费Kafka的消息、定时任务、普通HTTP接口都混在同一个gRPC线程池里。最好给不同业务分配独立的线程池,就像餐厅有专门的传菜员和服务员,各管一摊,一个环节出问题不会拖垮所有环节。

5.4 故障演练与监控告警

咱们得把“线程池耗尽”当成一个可能发生的故障来演练。比如故意把下游服务调慢,观察gRPC线程池活跃数量、队列深度、拒绝次数、响应时间等指标。监控指标至少要包括下面几个:

  • 活跃线程数
  • 队列中任务数
  • 拒绝任务次数
  • 线程池任务完成速率
  • 请求平均耗时和P99耗时
  • 下游服务的超时时间和超时比例

有了这些指标,咱们才能尽早发现隐患。比如活跃线程数长期超过80%,或者队列长度持续增长,就需要马上排查是不是下游慢了。不要等到线程池满了才报警,那已经晚了。

再看一个完整的超时配置示例。gRPC客户端可以设置超时时间,避免无限等待。

// 文件: GrpcClientMain.java
package demo;

import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.CallOptions;
import io.grpc.ClientCall;
import io.grpc.Metadata;
import io.grpc.Status;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.TimeUnit;

public class GrpcClientMain {

    public static void main(String[] args) throws Exception {
        ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 9000)
                .usePlaintext()
                .build();

        GreeterGrpc.GreeterBlockingStub blockingStub = GreeterGrpc.newBlockingStub(channel);

        // 构建请求
        demo.HelloRequest request = demo.HelloRequest.newBuilder()
                .setName("张三")
                .build();

        // 设置2秒超时,防止服务端一直不返回
        blockingStub = blockingStub.withDeadlineAfter(2, TimeUnit.SECONDS);

        try {
            demo.HelloReply reply = blockingStub.sayHello(request);
            System.out.println("收到响应: " + reply.getMessage());
        } catch (Exception e) {
            // 打印异常,一般会是DEADLINE_EXCEEDED
            System.out.println("调用失败: " + Status.fromThrowable(e));
        } finally {
            channel.shutdownNow();
        }
    }
}

5.5 代码实现综合示例

下面我们把前面的思路整合起来,写一个更完整的服务端:使用有界线程池、设置合理的参数、并用异步Stub和超时控制。技术栈还是Java。这个示例不是一个可以一键运行的完整项目,但主要部分都有了。

// 文件: GrpcServerWithGovernance.java
package demo;

import io.grpc.Server;
import io.grpc.ServerBuilder;
import io.grpc.stub.ServerCallStreamObserver;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.*;

public class GrpcServerWithGovernance {

    public static void main(String[] args) throws Exception {
        // 创建一个精心设计的线程池:核心线程16,最大32,队列容量500,拒绝策略为“调用线程执行”
        ThreadPoolExecutor executor = new ThreadPoolExecutor(
                16,
                32,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(500),
                r -> {
                    Thread t = new Thread(r, "grpc-governed-worker");
                    t.setDaemon(true);
                    return t;
                },
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        // 注册服务,服务实现里会用到这个线程池
        Server server = ServerBuilder.forPort(9000)
                .addService(new GovernedGreeterService(executor))
                .executor(executor)
                .build()
                .start();

        System.out.println("治理后的gRPC服务已启动");
        server.awaitTermination();
    }
}

// 服务实现类
class GovernedGreeterService extends GreeterGrpc.GreeterImplBase {

    private final ThreadPoolExecutor businessExecutor;

    public GovernedGreeterService(ThreadPoolExecutor businessExecutor) {
        this.businessExecutor = businessExecutor;
    }

    @Override
    public void sayHello(HelloRequest request, StreamObserver<HelloReply> responseObserver) {
        // 明确告知客户端,本服务会支持取消
        if (responseObserver instanceof ServerCallStreamObserver) {
            ((ServerCallStreamObserver<HelloReply>) responseObserver).setOnCancelHandler(() -> {
                // 当客户端取消时,可以做一些清理工作(这里只是打印日志)
                System.out.println("客户端取消了请求");
            });
        }

        // 提交到业务线程池,而不是直接用gRPC的线程池
        businessExecutor.execute(() -> {
            try {
                // 模拟耗时操作
                Thread.sleep(300);
                HelloReply reply = HelloReply.newBuilder()
                        .setMessage("你好, " + request.getName())
                        .build();
                responseObserver.onNext(reply);
                responseObserver.onCompleted();
            } catch (Exception e) {
                responseObserver.onError(e);
            }
        });
    }
}

在这个综合示例里,我们使用了有界队列、自定义线程名、合理的拒绝策略、客户端取消回调。虽然还没有达到生产级那么复杂,但已经能应对不少场景了。

六、应用场景与优缺点

gRPC服务端线程池治理特别适用于以下场景:高并发、低延迟要求高的微服务调用;多依赖外部系统且外部系统性能不稳定的服务;需要快速失败、避免线程堆积的业务;以及大规模分布式系统中的关键链路。

同步阻塞模型的优点是代码简单、容易理解,特别适合业务逻辑简单的场景。缺点是线程利用率低,一旦下游变慢,线程池就容易被打满,而且故障会传染。异步非阻塞模型的优点是线程利用率高,能支撑更高的并发,缺点是代码复杂度高,排查问题也更难,需要比较好的编程功底。

线程池调参本身也有优缺点。调大线程池可以暂时缓解线程不够用,但如果任务都是阻塞型的,调大只会延长每个请求的排队时间,甚至提高上下文切换开销。调小线程池虽然能控制并发,但如果业务本身需要很高的并发吞吐,可能造成频繁排队。有界队列是个好选择,但队列长度设置也要结合业务实际情况,太长会让请求长时间等待,太短会频繁触发拒绝策略。

因此,没有一套配置是万能的,需要根据业务场景、下游情况、机器配置不断调整。调优的过程最好用压测来验证,比如用 ghz 或者 grpcurl 模拟不同并发量,观察线程池指标。

七、注意事项

第一,不要试图用“调大线程池”来解决所有问题。阻塞调用导致的线程堆积,线程再多也不够。要把阻塞的调用改成异步或者非阻塞。第二,线程池关闭时要优雅。如果服务正在重启,要让线程池先停止接受新任务,然后等待已提交的任务完成。否则用户请求可能会被中途切断。第三,注意线程池的监控。只配置不监控,等于没有治理。至少要记录线程活跃数、任务排队数、拒绝数。第四,gRPC客户端重试要谨慎。如果服务端已经过载,客户端重试会加重服务端压力。可以限制重试次数,或者使用指数退避策略。第五,链路追踪很重要。要在每个gRPC请求里带上traceId,当出现线程池耗尽时,能快速定位是哪一条调用链引起的。第六,不要忽略业务代码里的锁竞争。如果很多线程都在等待一把锁,那和等外部接口没什么区别,同样会占住线程池。第七,信号量限流要放在业务处理之前,而且要在 finally 里释放许可,否则一旦异常,许可就永远少了。

八、文章总结

gRPC服务端线程池耗尽并不是一个孤立的技术问题,它往往是由“同步阻塞调用 + 无界队列或超大线程池 + 缺乏流量控制 + 缺少监控预警”共同作用的结果。咱们从一次事故入手,讲解了线程池的工作机制,分析了阻塞调用如何一步步堵死线程池,还提到了操作系统在大量线程调度下的开销。更重要的是,咱们给出了一套完整的治理策略:用异步替代阻塞、合理配置线程池参数、引入信号量限流、做好线程池和调用的监控告警,最后还给出了综合的代码示例。

在实际工作中,大家一定要记住:线程池不是越大越好,阻塞调用是万恶之源。遇到类似故障时,先问自己几个问题:我的线程在等什么?能不能不等?等的时候别人能不能帮忙做?如果不能,我能不能限制数量?把这些想明白,gRPC服务就会稳得多。希望这篇文章对你有用,解决一些真正让你头疼的线上问题。