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

基于RxJS/Observable重写服务队列串行执行逻辑的需求

Rewrite Promise-based Sequential Service Queue with RxJS Observables

Let's break down how to convert your Promise-based sequential service execution into RxJS, preserving all core behavior: sequential execution, passing previous results to the next service, stopping immediately on failure, and handling both synchronous and asynchronous services.

First, let's recap your existing setup for context:

Original Promise Implementation

You have a collection of service definitions, and you use reduce with Promises to chain them sequentially:

let collection = [ { name: 'calcOnServer', parameters: {} }, { name: 'calc', parameters: {} } ];
return collection.reduce((currentFunction, nextFunction, index) => {
  return currentFunction.then(() => {
    let result = runFunction(nextFunction.name, nextFunction.parameters); // runFunction fetches the class and calls .run()
    if (result === undefined) {
      result = Promise.resolve();
    }
    if (!result.then) {
      if (Boolean(result)) {
        result = Promise.resolve(result);
      } else {
        result = Promise.reject();
      }
    }
    return result.then((result) => {
      collection[index + 1].parameters = result;
    });
  });
}, Promise.resolve())

Service Examples

Your services can be asynchronous (returning Promises) or synchronous (returning values or rejecting):

class calcOnServer {
  run({param1, param2}) {
    return new Promise((resolve, reject) => {
      // Simulate async operation
      setTimeout(() => {
        resolve({param1: 5, param2: 3});
      }, 100);
    });
  }
}

class calc {
  run({param1, param2}) {
    if (typeof param1 === 'number' && typeof param2 === 'number') {
      return param1 + param2;
    } else {
      return Promise.reject(new Error('Invalid parameters'));
    }
  }
}

// Assume runFunction is defined like this:
function runFunction(serviceName, params) {
  const service = new window[serviceName]();
  return service.run(params);
}

RxJS Solution

RxJS has built-in operators that fit this sequential execution pattern perfectly. We'll use from to convert your service array into an Observable, concatMap to enforce sequential execution, and handle result normalization to match your original Promise logic.

Here's the rewritten code:

import { from, of, throwError } from 'rxjs';
import { concatMap, tap, catchError } from 'rxjs/operators';

function executeServiceQueue(collection) {
  // Start processing the service collection as an Observable stream
  return from(collection).pipe(
    concatMap((serviceDef, index) => {
      // Execute the current service, same as your original logic
      let result = runFunction(serviceDef.name, serviceDef.parameters);

      // Normalize the result to an Observable, matching your Promise handling rules
      if (result === undefined) {
        result = of(undefined); // Equivalent to Promise.resolve()
      } else if (!(result.subscribe || result.then)) {
        // Handle synchronous return values: resolve if truthy, reject if falsy
        result = Boolean(result) ? of(result) : throwError(() => new Error('Service returned falsy value'));
      } else if (result.then) {
        // Convert Promises to Observables
        result = from(result);
      }

      // Update the next service's parameters (if it exists) after successful execution
      return result.pipe(
        tap((resultValue) => {
          if (collection[index + 1]) {
            collection[index + 1].parameters = resultValue;
          }
        })
      );
    }),
    // Optional: Catch and log errors, then re-throw to maintain failure termination behavior
    catchError((err) => {
      console.error('Service queue failed:', err);
      return throwError(() => err);
    })
  );
}

// Usage example:
const collection = [ { name: 'calcOnServer', parameters: {} }, { name: 'calc', parameters: {} } ];
executeServiceQueue(collection).subscribe({
  next: (stepResult) => console.log('Step completed with result:', stepResult),
  error: (err) => console.error('Queue terminated due to error:', err),
  complete: () => console.log('All services executed successfully!')
});

Key Details Explained:

  • from(collection): Converts your service array into an Observable that emits each service definition one at a time.
  • concatMap: Guarantees sequential execution—we wait for the current service's Observable to complete before moving to the next one, just like your Promise chain.
  • Result Normalization: We replicate every edge case from your original code:
    • Undefined results are treated as successful (empty resolution)
    • Synchronous truthy values are wrapped as successful Observables
    • Synchronous falsy values trigger an error
    • Promises are converted to Observables for consistency
  • tap: Used to update the next service's parameters after a successful step, mirroring your original .then() mutation of the collection array.
  • Error Handling: Any error (rejected Promise, thrown error, or falsy sync value) immediately terminates the stream, stopping all subsequent services—exactly like your original Promise implementation.

Immutable Alternative (No Array Mutation)

If you want to avoid mutating the original collection array (to prevent side effects), you can use scan to carry the result forward to the next service instead:

import { from, of, throwError } from 'rxjs';
import { concatMap, scan, catchError } from 'rxjs/operators';

function executeServiceQueueImmutable(collection) {
  return from(collection).pipe(
    scan((previousResult, serviceDef) => {
      // Use the previous result as parameters for the current service (fallback to original params if first step)
      const params = previousResult ?? serviceDef.parameters;
      let result = runFunction(serviceDef.name, params);

      // Normalize result to Observable
      if (result === undefined) result = of(undefined);
      else if (!(result.subscribe || result.then)) result = Boolean(result) ? of(result) : throwError(() => new Error('Falsy value'));
      else if (result.then) result = from(result);

      return result;
    }, undefined),
    concatMap(obs => obs), // Flatten the Observable of Observables
    catchError(err => throwError(() => err))
  );
}

This version keeps the original collection intact and passes results directly between services via the scan operator.


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:32:03