如何在NestJS中用WebSocket和RxJs实现长运行进程的Observable监控
The issue in your code is that you're passing the longRunningProcess function itself to from() instead of invoking it to get the Observable it returns. Plus, you need to map each value emitted by the process to the WsResponse format that NestJS WebSocket expects.
Here's how to fix your findAll method:
@SubscribeMessage('events') findAll(@MessageBody() data: any): Observable<WsResponse<unknown>> { return longRunningProcess().pipe( map((message) => ({ event: 'events', data: message })) ); }
Breakdown of the fix:
- Invoke the function:
longRunningProcess()returns the Observable that emits your step updates, so we use this directly instead of wrapping it infrom().from()is for converting arrays/promises/iterables to Observables, but since your function already returns an Observable, we don't need it here. - Map to WsResponse: The
pipe(map(...))transforms each emitted value (like "End of step 1") into the structure NestJS requires for WebSocket responses. This matches the working example you had withfrom([1,2,3]).
How it works:
When a client sends the events message, the server will subscribe to the Observable returned by longRunningProcess(). Each time subscriber.next() is called in your long-running process, the mapped WsResponse will be sent over the WebSocket to the client. Finally, when subscriber.complete() is triggered after the timeout, the Observable completes cleanly.
Expected client flow:
After the fix, your client will receive these messages in sequence:
{"event":"events","data":"End of step 1"}{"event":"events","data":"End of step 2"}{"event":"events","data":"End of step 3"}{"event":"events","data":"End of Step 4"}
内容的提问来源于stack exchange,提问作者Martin Thompson

