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

如何基于manualForceStop让Poll Observable主动抛错并停止轮询

手动终止轮询并抛出指定错误的实现方案

当前你的poll函数通过retry实现轮询逻辑,但无法响应manualForceStop的终止信号——即使将其设为true,retry仍会持续重试。要实现手动终止并抛出MANUAL_FORCE_STOP错误,需要把终止信号融入到轮询的整个生命周期中,以下是两种可行方案:

方案一:基于retry的改进实现

通过在retry的延迟逻辑中加入终止信号监听,同时在全局流中捕获终止信号,确保触发终止时立即停止重试并抛出错误:

const CONVERT_SECONDS_TO_MS = 1000;
interface PollCmd<TApiArgs, TApiResponse> {
  intervalInSeconds: number;
  timeoutInSeconds: number;
  apiCall: (arg?: TApiArgs) => Observable<TApiResponse>;
  shouldContinuePolling: (arg: TApiResponse) => boolean;
  manualForceStop: BehaviorSubject<boolean>;
}

export function poll<TApiArgs, TApiResponse>(
  input: PollCmd<TApiArgs, TApiResponse>
): Observable<TApiResponse> {
  // 构建手动终止信号流:当manualForceStop变为true时抛出指定错误
  const manualStop$ = input.manualForceStop.pipe(
    filter(stop => stop),
    switchMap(() => throwError(() => 'MANUAL_FORCE_STOP'))
  );

  return input.apiCall().pipe(
    switchMap((apiResponse) => {
      if (input.shouldContinuePolling(apiResponse)) {
        return throwError(() => 'CONDITION_NOT_REACHED');
      }
      return of(apiResponse);
    }),
    // 自定义重试逻辑,加入终止判断
    retry({
      delay: (error) => {
        // 若为手动终止错误,直接抛出不再重试
        if (error === 'MANUAL_FORCE_STOP') {
          return throwError(() => error);
        }
        // 等待重试间隔的同时监听终止信号
        return timer(input.intervalInSeconds * CONVERT_SECONDS_TO_MS).pipe(
          takeUntil(manualStop$)
        );
      }
    }),
    timeout({ first: input.timeoutInSeconds * CONVERT_SECONDS_TO_MS }),
    // 全局监听终止信号,确保API调用过程中也能响应终止
    takeUntil(manualStop$),
    // 确保手动终止错误正确向上抛出
    catchError((error) => {
      if (error === 'MANUAL_FORCE_STOP') {
        return throwError(() => error);
      }
      throw error;
    })
  );
}

关键逻辑说明

  1. manualStop$流:过滤出manualForceStop为true的信号,将其转换为抛出指定错误的流,作为终止触发源。
  2. retry延迟处理:在重试等待间隔中通过takeUntil监听终止信号,一旦触发就停止等待并抛出错误;同时判断错误类型,避免对终止错误进行重试。
  3. 全局takeUntil:确保即使在API调用执行过程中触发终止,也能立即终止整个流。

方案二:基于repeatWhen的更可控实现

repeatWhen比retry更灵活,可以精准控制哪些错误需要触发重试,同时融入终止信号:

const CONVERT_SECONDS_TO_MS = 1000;
interface PollCmd<TApiArgs, TApiResponse> {
  intervalInSeconds: number;
  timeoutInSeconds: number;
  apiCall: (arg?: TApiArgs) => Observable<TApiResponse>;
  shouldContinuePolling: (arg: TApiResponse) => boolean;
  manualForceStop: BehaviorSubject<boolean>;
}

export function poll<TApiArgs, TApiResponse>(
  input: PollCmd<TApiArgs, TApiResponse>
): Observable<TApiResponse> {
  const manualStop$ = input.manualForceStop.pipe(
    filter(stop => stop),
    tap(() => { throw new Error('MANUAL_FORCE_STOP'); })
  );

  return input.apiCall().pipe(
    switchMap(apiResponse => {
      if (input.shouldContinuePolling(apiResponse)) {
        return throwError(() => 'CONDITION_NOT_REACHED');
      }
      return of(apiResponse);
    }),
    // 仅对CONDITION_NOT_REACHED错误触发重试
    repeatWhen(errors => 
      errors.pipe(
        filter(err => err === 'CONDITION_NOT_REACHED'),
        mergeMap(() => 
          timer(input.intervalInSeconds * CONVERT_SECONDS_TO_MS).pipe(
            takeUntil(manualStop$)
          )
        )
      )
    ),
    timeout({ first: input.timeoutInSeconds * CONVERT_SECONDS_TO_MS }),
    takeUntil(manualStop$)
  );
}

关键逻辑说明

  1. repeatWhen过滤错误:只处理CONDITION_NOT_REACHED错误,其他错误(包括手动终止错误)会直接终止重试流程。
  2. 终止信号触发:manualStop$中的tap会直接抛出错误,当takeUntil触发时,整个流会立即终止并抛出该错误。
  3. 重试间隔监听:在等待重试的timer流中加入takeUntil,确保终止信号能打断重试等待。

两种方案都能满足需求,你可以根据自己的代码风格选择。核心思路是让终止信号能渗透到轮询的每一个环节——无论是API调用中、重试等待中还是重试判断中,确保触发终止后立即停止所有操作并抛出指定错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:59:57