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

如何重置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 the resetTimer$.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 retryWhen instead of blindly resetting the timer.
  • Side Effect Best Practices: Use tap instead of map for writing to the SSE response — map should only transform data, not execute side effects.

内容的提问来源于stack exchange,提问作者user19639900

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 20:59:10