如何重置RxJS定时器?以及当出现401授权错误时如何重置该定时器?
Alright, let's break this down step by step. First, I'll cover how to reset an RxJS timer in general, then we'll adapt that logic to your SSE endpoint where you need to handle 401 unauthorized errors.
1. How to Reset an RxJS Timer (General Approach)
RxJS's built-in timer is a cold observable — once it starts, it runs on its fixed interval with no built-in way to reset it. To add reset functionality, we use a Subject as a "reset trigger" combined with switchMap, which cancels the previous subscription whenever a new signal comes in.
Here's a simple example:
import { Subject, timer, switchMap, startWith } from 'rxjs'; // Create a Subject to act as our reset trigger const resetTimer$ = new Subject<void>(); // Build a resettable timer stream const resettableTimer$ = resetTimer$.pipe( startWith(undefined), // Trigger the initial timer on startup switchMap(() => timer(0, 44000)) // Cancel old timer, start new one on reset signal ); // Call this function whenever you need to reset the timer function triggerTimerReset() { resetTimer$.next(); }
When triggerTimerReset() is called, switchMap will terminate the current timer subscription and spin up a fresh one, effectively resetting the interval.
2. Adapting to Your SSE Endpoint & 401 Error Handling
Looking at your existing code, we need to:
- Wrap your async
postService.count()call in observable logic to catch 401 errors - Integrate the reset mechanism to restart the timer when a 401 occurs
- Clean up the SSE response handling to follow RxJS best practices
Here's the modified code with explanations:
import { Subject, timer, switchMap, tap, catchError, EMPTY } from 'rxjs'; import { HttpException } from '@nestjs/common'; @Sse('true') @HasRole(Role.ADMIN) async sse(@Res() response: Response) { // 1. Set up SSE response headers first response.setHeader('Content-Type', 'text/event-stream'); response.setHeader('Cache-Control', 'no-cache'); response.setHeader('Connection', 'keep-alive'); // 2. Create our reset trigger Subject const resetTimer$ = new Subject<void>(); // 3. Build the resettable SSE stream const sseStream$ = resetTimer$.pipe( startWith(undefined), // Start the timer immediately on connection switchMap(() => timer(0, 44000).pipe( switchMap(async () => { try { return await this.postService.count(); } catch (err) { // Catch 401 errors specifically if (err instanceof HttpException && err.getStatus() === 401) { console.log('401 Unauthorized detected — resetting timer'); // Trigger the reset to restart the timer/request cycle resetTimer$.next(); } throw err; // Re-throw to let catchError handle other errors } }), // Write SSE data as a side effect (use tap instead of map for side effects) tap(result => { response.write(`data: ${result.count}\n\n`); }), catchError(err => { console.error('Error in SSE stream:', err); // Optional: Add retry limits here to avoid infinite loops // For now, reset the timer to keep the stream alive resetTimer$.next(); return EMPTY; // Continue the stream after error }) ) ) ); // 4. Subscribe to the stream to keep the SSE connection open sseStream$.subscribe({ error: () => response.end(), // Close connection on unrecoverable errors complete: () => response.end() }); return sseStream$; }
Key Notes for Your Use Case:
- 401 Handling: When a 401 is caught from
postService.count(), we trigger the timer reset immediately. If your flow requires waiting for a new auth token (e.g., after user re-login), you could tie theresetTimer$.next()call to an auth state change event instead. - Avoid Infinite Loops: If the 401 persists (e.g., invalid credentials), add a retry limit using
retryWheninstead of blindly resetting the timer. - Side Effect Best Practices: Use
tapinstead ofmapfor writing to the SSE response —mapshould only transform data, not execute side effects.
内容的提问来源于stack exchange,提问作者user19639900

