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
相关产品推荐
相关产品推荐

