ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RxJS 4 异步生成器实战:`Rx.Observable.spawn` 完整解析与源码级原理

RxJS 4 异步生成器实战:`Rx.Observable.spawn` 完整解析与源码级原理 RxJS 4 异步生成器实战Rx.Observable.spawn完整解析与源码级原理【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS导读Rx.Observable.spawn是 RxJS v4本仓库即 Reactive-Extensions/RxJS 的 v4 代码库中一个极具特色的静态工厂方法它允许你用一个 ES6 生成器函数Generator Function以近乎同步的书写风格依次yield等待 Promise、Observable、数组、对象、Node 风格回调函数乃至嵌套生成器的结果最终产出一个携带最终返回值的Observable。阅读本文后你将掌握spawn的完整 API 语义、各类可等待值的转换规则、错误处理机制、配套的Observable.wrap用法并能结合 spawn 源码 与 单元测试 理解其底层驱动循环的实现原理直接在你的异步链式代码中落地使用。一、spawn是什么用生成器编排异步序列在 RxJS v4 中异步编排通常依赖操作符的层层嵌套flatMap、concatMap等。spawn换了一种思路把生成器函数作为剧本yield出来的每一个值都被当作一个待等待的异步任务任务完成后其结果会被注入回生成器继续往下执行直到生成器return出最终结果该结果会作为Observable的唯一onNext值发射并随即onCompleted。它在 src/core/linq/observable/spawn.js 中实现官方 API 文档位于 doc/api/core/operators/spawn.md。本质上它把回调地狱或操作符嵌套改写成了顺序直观的同步风格代码与后来 ES2017 的async/await思路一致——区别在于它等待的不只是 Promise还包括 Observable、thunk、数组、对象与生成器。二、API 签名参数与返回值spawn是Rx.Observable上的静态方法其签名定义在 TypeScript 声明 ts/core/linq/observable/spawn.ts 中module Rx { export interface ObservableStatic { wrapT(fn: Function): ObservableT; spawnT(fn: Function): ObservableT; } }参数fnFunction生成器函数function* () { ... }也可以是普通函数——此时函数会被先调用一次若其返回值为生成器对象gen.next是函数则继续按生成器流程处理否则直接把该返回值作为结果发射。返回Observable一个携带最终结果的Observable。生成器return的值作为onNext发射随后onCompleted中途任何异常或异步任务失败则走onError。除spawn外同文件还一并定义了Observable.wrap见第五节。三、官方示例一次等待五种不同类型原文档给出的示例完整展示了spawn的核心用法在同一个生成器里依次yield一个 Node 风格回调thunk、一个数组、两个 Observable 和一个 Promise然后像同步代码一样把它们拼接成最终结果var Rx require(rx); var spawned Rx.Observable.spawn(function* () { var a yield cb cb(null, a); var b yield [b]; var c yield Rx.Observable.just(c); var d yield Rx.Observable.just(d); var e yield Promise.resolve(e); return a b c d e; }); spawned.subscribe( function (x) { console.log(next %s, x); }, function (e) { console.log(error %s, e); }, function () { console.log(completed); } ); // next abcde // completed执行结果a拿到回调值ab拿到数组[b]c/d拿到两个Observable.just的值e拿到 Promise 的解析值最终return的拼接结果abcde作为next发射随后立即completed。这段代码可以在 Node 环境require(rx)直接运行或者引入本仓库构建产物后运行。四、yield 可等待值全解析toObservable的转换规则spawn最强大的地方在于它认识的可等待值类型非常丰富。核心转换逻辑集中在 spawn.js 中的toObservablefunction toObservable(obj) { if (!obj) { return obj; } if (Observable.isObservable(obj)) { return obj; } if (isPromise(obj)) { return Observable.fromPromise(obj); } if (isGeneratorFunction(obj) || isGenerator(obj)) { return spawn.call(this, obj); } if (isFunction(obj)) { return thunkToObservable.call(this, obj); } if (isArrayLike(obj) || isIterable(obj)) { return arrayToObservable.call(this, obj); } if (isObject(obj)) {return objectToObservable.call(this, obj);} return obj; }对每种类型的处理细节如下1. Observable原样保留Observable.isObservable(obj)为真时直接返回不做任何包装。yield Rx.Observable.just(c)等待的就是该序列发射的最后一个值。2. Promise自动转 Observable通过Observable.fromPromise包装。Promise 解析出的值注入生成器Promise 拒绝时走onError分支。3. Node 风格回调函数 / thunkthunkToObservable对应官方示例里的cb cb(null, a)。所谓 thunk 是一个只接受回调作为唯一参数的函数其回调遵循 Node 惯例(err, result)。实现见 thunkToObservablefunction thunkToObservable(fn) { var self this; return new AnonymousObservable(function (o) { fn.call(self, function () { var err arguments[0], res arguments[1]; if (err) { return o.onError(err); } if (arguments.length 2) { var args []; for (var i 1, len arguments.length; i len; i) { args.push(arguments[i]); } res args; } o.onNext(res); o.onCompleted(); }); }); }规则很清晰回调的第一个参数非空即视为错误并onError回调携带超过 2 个参数时多余参数会被收集成数组作为结果。这意味着你可以在spawn里直接等待fs.readFile这类 Node API 的 thunk 化版本。4. 数组 / 类数组 / 可迭代对象arrayToObservableyield [b]走 arrayToObservablefunction arrayToObservable (obj) { return Observable.from(obj).concatMap(function(o) { if(Observable.isObservable(o) || isObject(o)) { return toObservable.call(null, o); } else { return Rx.Observable.just(o); } }).toArray(); }它把数组里的每个元素串行concatMap等待元素若是 Observable 或对象则继续递归转换否则直接just包装最后toArray()把整组结果聚合成一个新数组注入生成器。这正是官方示例里b拿到[b]的原因。5. 普通对象并行forkJoinyield 一个对象 时所有属性值中的 Observable 会被并行等待然后合并回原对象结构function objectToObservable (obj) { var results new obj.constructor(), keys Object.keys(obj), observables []; for (var i 0, len keys.length; i len; i) { var key keys[i]; var observable toObservable.call(this, obj[key]); if(observable Observable.isObservable(observable)) { defer(observable, key); } else { results[key] obj[key]; } } return Observable.forkJoin.apply(Observable, observables).map(function() { return results; }); }非 Observable 属性原样保留Observable 属性通过forkJoin并行等待全部完成后用map把填充好的results对象返回。注意new obj.constructor()意味着原对象必须是可构造的例如普通对象字面量返回的是同一类型的新对象。6. 生成器 / 生成器函数嵌套spawnyield出一个生成器或生成器函数时会递归调用spawn把它当作一个内层异步序列来等待支持任意层级的嵌套组合。7. 不支持的原始值toObservable返回非 Observable 值如数字、字符串、null等时spawn 主循环会抛出TypeError(type not supported)并走onError——也就是说你不能直接yield 42这种裸值需要包一层Observable.just(42)。这一点可以从 spawn 主循环的next函数 得到印证function next(ret) { if (ret.done) { o.onNext(ret.value); o.onCompleted(); return; } var obs toObservable.call(self, ret.value); var value null; var hasValue false; if (Observable.isObservable(obs)) { g.add(obs.subscribe(function(val) { hasValue true; value val; }, onError, function() { hasValue processGenerator(value); })); } else { onError(new TypeError(type not supported)); } }五、源码级原理spawn 的驱动循环spawn的完整实现位于 src/core/linq/observable/spawn.js其核心是一个AnonymousObservable内部通过三个函数接力驱动生成器processGenerator(res)把上一个异步任务的结果res通过gen.next(res)重新注入生成器用tryCatch捕获生成器内部同步抛出的异常异常直接o.onError拿到{ value, done }后交给next。onError(err)异步任务失败时调用gen.next(err)注意是next而非throw把错误值作为普通值注入生成器——这意味着你可以在生成器内部用try/catch包裹yield来捕获异步错误实现类似async/await的错误处理体验。若生成器本身不处理该错误异常最终会经onError传播给订阅者。next(ret)若ret.done为真将ret.value作为onNext发射并onCompleted否则把ret.value转成 Observable 订阅之订阅成功拿到值后用processGenerator继续推进。整个流程用CompositeDisposableg管理当前待处理的订阅保证 Observable 被取消订阅时中间任务能一并清理。整个循环从processGenerator()无参启动开始直至生成器done或出错终止。一个值得注意的设计细节if (isFunction(gen)) { gen gen.apply(self, args); }支持给生成器函数传参spawn(fn, arg1, arg2, ...)也允许传入返回生成器对象的普通函数。六、Observable.wrap把生成器封装成可复用的 Observable 工厂与spawn配套的Observable.wrap定义在同一文件的 前 8 行Observable.wrap function (fn) { function createObservable() { return Observable.spawn.call(this, fn.apply(this, arguments)); } createObservable.__generatorFunction__ fn; return createObservable; };它接收一个生成器函数返回一个新的函数每次调用该函数都会先用当前调用参数调用原生成器函数拿到生成器实例再交给spawn生成 Observable。这相当于spawn的惰性化/参数化变体你可以在多处、以不同参数反复调用同一个包装函数各自产生独立的 Observable 序列非常适合把一段异步流程固化成可复用工厂var fetchFlow Rx.Observable.wrap(function* (url) { var data yield getJsonThunk(url); // thunk var enriched yield Rx.Observable.fromPromise(loadMore(data)); return enriched; }); fetchFlow(/api/a).subscribe(...); fetchFlow(/api/b).subscribe(...); // 再次调用互不影响七、错误处理与完成语义结合主循环代码可以总结出spawn的完整错误与完成语义场景行为源码依据生成器return值onNext(value)后立即onCompletedspawn.js#L38-L42生成器内部同步抛异常tryCatch捕获直接onErrorspawn.js#L23-L25等待的 Promise 拒绝 / Observable 出错 / thunk 回调带 err通过gen.next(err)注入生成器可在内部try/catch捕获spawn.js#L31-L35yield 了不支持的类型TypeError(type not supported)走onErrorspawn.js#L53-L55传入的fn不是生成器且无next方法直接把gen作为唯一值发射并完成spawn.js#L17-L21取消订阅通过CompositeDisposable释放中间订阅spawn.js#L15八、测试与类型定义验证仓库在 tests/observable/spawn.js 中提供了 QUnit 测试用例QUnit.module(spawn)。其中spawn success 1用例用TestScheduler虚拟时间驱动生成器first依次yield两个Rx.Observable.just42 与 56return x y最终断言在时刻 200 收到onNext(200, 98)与onCompleted(200)见 测试第 698-714 行。该测试验证了 spawn 的基本成功路径串行等待两个 Observable、正确求和并把结果作为唯一 next 发射后完成。由于测试文件中生成器经过 regenerator 编译可直接阅读其编译形态理解执行流程。类型层面spawn与wrap均声明在 ts/core/linq/observable/spawn.ts 的ObservableStatic接口中并随 ts/rx.all.d.ts 等聚合声明文件一起对外发布表明两者属于官方公开 API 而非内部实现细节。九、如何获取与使用spawn属于 RxJS 的async 能力范畴随各发行形态发布。原文档doc/api/core/operators/spawn.md指出其位于rx.async/rx.all系列构建中并依赖基础rx或rx.lite与rx.binding模块NPM 包为rxNuGet 包为RxJS-All与RxJS-Async。本仓库对应的发布形态如下模块化源码实现见 modules/rx-lite-async/rx.lite.async.js含spawn/wrap/toObservable及全部辅助函数以及压缩版 modules/rx-lite-async/rx.lite.async.min.js对应包配置见 modules/rx-lite-async/package.json聚合构建rx.all系列RxJS-All与rx.async系列RxJS-Async的 NuGet 规格见 nuget/RxJS-All/RxJS-All.nuspec 与 nuget/RxJS-Async/RxJS-Async.nuspec核心源码可读性最佳始终建议直接阅读 src/core/linq/observable/spawn.js该文件同时定义了spawn、wrap与全部转换辅助函数注释与结构最清晰。在 Node 中使用时v4 版本需支持生成器的运行时或通过 regenerator 转译引入rx包后即可按第三节示例直接调用Rx.Observable.spawn。注意本仓库为 RxJS v4 代码库文档开头即标注 This is RxJS v 4与 RxJS 5/6/7 的 API 并不兼容。十、综合实战把多源异步流程写成同步代码最后用一个综合示例收尾把本文讲到的所有要点串起来一个生成器依次等待 thunkNode 风格回调、Promise、Observable、嵌套生成器和一个对象并行等待多个 Observable 属性最后汇总返回var Rx require(rx); function readFileThunk(path) { return function (cb) { // 模拟异步读取 setTimeout(function () { cb(null, path :content); }, 10); }; } function* innerGen(x) { var y yield Rx.Observable.just(x * 2); return y 1; } var flow Rx.Observable.spawn(function* (path) { var file yield readFileThunk(path); // 等待 thunk var fromPromise yield Promise.resolve(P); // 等待 Promise var fromObs yield Rx.Observable.just(O); // 等待 Observable var nested yield innerGen(3); // 等待嵌套生成器 - 7 var obj yield { a: Rx.Observable.just(A), b: Promise.resolve(B), plain: raw // 非 Observable 属性原样保留 }; // 并行等待返回 {a:A, b:B, plain:raw} return [file, fromPromise, fromObs, nested, obj]; }); flow(README).subscribe( function (x) { console.log(next:, JSON.stringify(x)); }, function (e) { console.log(error:, e); }, function () { console.log(completed); } ); // next: [README:content,P,O,7,{a:A,b:B,plain:raw}] // completed可以看到无论底层数据源是回调、Promise、Observable 还是并行对象spawn都把它们统一成了同步顺序的阅读体验同时保留了 Observable 的完整错误传播与取消订阅能力。理解其toObservable转换表和驱动循环之后你完全可以基于 src/core/linq/observable/spawn.js 的范式在自己的项目中实现类似的生成器式异步编排工具。【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表