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

Nest拦截器导致Azure Service Bus集成故障排查求助

问题分析与解决方案

核心原因

你的自定义Azure Service Bus传输策略未正确处理Nest拦截器返回的Observable流。添加拦截器后,处理器方法的返回值会被包装成Observable,但listen方法中直接用await handler(...)等待该Observable——await Observable不会等待它发出实际值,而是直接返回Observable实例本身。这个实例被序列化成JSON后,就出现了{"source":{"source":{}}}的异常响应。

同时,拦截器中直接修改原始请求对象(delete message.sicil)存在引用副作用风险,可能导致后续逻辑拿到不完整的请求数据。

修复步骤

1. 修正自定义传输的listen方法

使用RxJS的lastValueFrom将Observable转换为Promise,确保拿到实际响应数据:

import { lastValueFrom } from 'rxjs';

// 其他代码保持不变
async listen(callback: () => void) {
  this.queueReceiver.subscribe({
    processMessage: async (brokeredMessage) => {
      const handler = this.getHandlerByPattern(brokeredMessage.sessionId);
      if (!handler) {
        console.error(`No handlers for pattern ${brokeredMessage.subject}`);
        return;
      } else {
        const result$ = handler(this.deserialize(brokeredMessage).data);
        // 将Observable转换为Promise,获取实际响应值
        const result = await lastValueFrom(result$);
        await this.queueSender.sendMessages({
          body: result,
          sessionId: brokeredMessage.sessionId,
        });
      }
    },
    processError: async (err) => {
      console.error(err);
    }
  });
  callback();
}

2. 避免修改原始请求对象

在拦截器中创建请求对象的副本,不直接修改原始数据,消除引用副作用:

intercept(context: ExecutionContext, next: CallHandler): Observable<any> {
  const eventName = context.getHandler().name;

  if (eventName === 'validate') {
    return next.handle();
  }

  const message = context.switchToRpc().getData();
  const sicil = message.id;
  // 创建副本,不修改原始对象
  const cleanedMessage = { ...message };
  delete cleanedMessage.sicil;

  // 若需要让控制器接收清理后的数据,添加以下代码
  // context.switchToRpc().setData(cleanedMessage);

  return next.handle().pipe(
    map((data) => {
      const payload: IncomingEventPayload = {
        eventName: eventName,
        data: data,
        event_date: new Date(),
        sicil: sicil,
      };
      this.client.emitEvent('auth-event', payload);
      return data;
    }),
    catchError((err) => {
      const payload: IncomingEventErrorPayload = {
        eventName: eventName,
        err: err,
        err_date: new Date(),
        sicil: sicil,
      };
      this.client.emitEvent('auth-error', payload);
      throw err;
    })
  );
}

3. 优化拦截器的异步操作处理

如果ClientProxyService.emitEvent是异步操作,建议用tap操作符替代map,避免阻塞主响应流:

// 替换原map操作符
tap(async (data) => {
  const payload: IncomingEventPayload = {
    eventName: eventName,
    data: data,
    event_date: new Date(),
    sicil: sicil,
  };
  await this.client.emitEvent('auth-event', payload);
}),

验证效果

修改后重启服务,添加拦截器测试,响应应恢复为正常格式:RESULT{"maskedGsm":"53******45"}

内容的提问来源于stack exchange,提问作者Emir Kutlugün

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:27:40