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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 14:17:59