基于RxJS/Observable重写服务队列串行执行逻辑的需求
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 thecollectionarray.- 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

