NestJS Worker消费GCP Pub/Sub:SIGTERM无法停止消费问题
解决NestJS Pub/Sub Worker SIGTERM无法停止消费的问题
1. 启用NestJS Shutdown Hooks
NestJS默认不会自动处理SIGTERM等终止信号,必须手动启用 shutdown 钩子,确保应用能捕获信号并执行清理逻辑:
// main.ts import { NestFactory } from '@nestjs/core'; import { AppModule } from './app.module'; async function bootstrap() { // 用ApplicationContext创建无HTTP服务的Worker实例 const app = await NestFactory.createApplicationContext(AppModule); // 启用 shutdown 钩子,捕获SIGTERM/SIGINT信号 app.enableShutdownHooks(); const pubsubService = app.get(PubsubService); await pubsubService.startConsuming(); // 保持进程运行(无HTTP服务时需要) process.stdin.resume(); } bootstrap();
2. 管理Pub/Sub订阅实例并实现优雅关闭
必须保存订阅实例的引用,在应用 shutdown 时主动关闭订阅,停止接收新消息。利用OnApplicationShutdown生命周期钩子触发清理:
// pubsub.service.ts import { Injectable, OnApplicationShutdown } from '@nestjs/common'; import { PubSub, Subscription } from '@google-cloud/pubsub'; @Injectable() export class PubsubService implements OnApplicationShutdown { private subscription: Subscription; constructor(private readonly pubsub: PubSub) {} async startConsuming() { this.subscription = this.pubsub.subscription('your-subscription-id'); // 绑定消息处理和错误监听 this.subscription.on('message', this.handleMessage.bind(this)); this.subscription.on('error', this.handleError.bind(this)); } private handleMessage(message) { // 你的消息处理逻辑 console.log(`处理消息: ${message.data.toString()}`); message.ack(); } private handleError(error) { console.error(`订阅错误: ${error}`); } // 应用终止时触发的钩子 async onApplicationShutdown(signal?: string) { if (this.subscription) { console.log(`收到信号${signal},正在停止Pub/Sub订阅...`); // 关闭订阅,停止接收新消息 await this.subscription.close(); console.log('Pub/Sub订阅已成功停止'); } } }
3. 优化CancellationTokenService的配合逻辑
如果用了取消令牌,要确保在shutdown时触发取消,并在消息处理中检查令牌状态,避免处理新消息:
// cancellationToken.service.ts import { Injectable } from '@nestjs/common'; @Injectable() export class CancellationTokenService { private isCancelled = false; cancel() { this.isCancelled = true; } get isCancellationRequested(): boolean { return this.isCancelled; } }
修改PubsubService,集成取消令牌:
// pubsub.service.ts // ... 其他导入和代码 constructor( private readonly pubsub: PubSub, private readonly cancellationTokenService: CancellationTokenService ) {} private handleMessage(message) { // 如果已触发取消,直接nack消息并返回 if (this.cancellationTokenService.isCancellationRequested) { message.nack(); return; } // 正常处理逻辑 console.log(`处理消息: ${message.data.toString()}`); message.ack(); } async onApplicationShutdown(signal?: string) { // 先触发取消令牌,阻止新消息处理 this.cancellationTokenService.cancel(); if (this.subscription) { console.log(`收到信号${signal},正在停止Pub/Sub订阅...`); await this.subscription.close(); console.log('Pub/Sub订阅已成功停止'); } }
4. 关键注意事项
- 确保
@google-cloud/pubsub依赖为最新版本,旧版本的close()方法可能存在兼容性问题。 - 如果消息处理包含异步操作,需在shutdown时等待当前正在处理的消息完成,或通过取消令牌中断未完成的操作。
- 测试时用
kill -SIGTERM <进程PID>发送终止信号,检查日志是否输出订阅停止信息,且进程能优雅退出。
内容的提问来源于stack exchange,提问作者Vinicius Andrade
相关产品推荐
相关产品推荐

