【发布时间】:2016-03-09 16:23:15
【问题描述】:
我正在使用 Rx.js 处理文件的内容,为每一行发出一个 http 请求,然后汇总结果。但是,源文件包含数千行,并且我正在重载执行 http 请求的远程 http api。我需要确保在开始另一个请求之前等待现有的 http 请求回调。我愿意一次批处理和执行n 请求,但是对于这个脚本来说,串行执行请求就足够了。
我有以下几点:
const fs = require('fs');
const rx = require('rx');
const rxNode = require('rx-node');
const doHttpRequest = rx.Observable.fromCallback((params, callback) => {
process.nextTick(() => {
callback('http response');
});
});
rxNode.fromReadableStream(fs.createReadStream('./source-file.txt'))
.flatMap(t => t.toString().split('\r\n'))
.take(5)
.concatMap(t => {
console.log('Submitting request');
return doHttpRequest(t);
})
.subscribe(results => {
console.log(results);
}, err => {
console.error('Error', err);
}, () => {
console.log('Completed');
});
但是,这不会以串行方式执行 http 请求。它输出:
提交请求 提交请求 提交请求 提交请求 提交请求 http响应 http响应 http响应 http响应 http响应 完全的如果我删除对 concatAll() 的调用,那么请求是串行的,但我的订阅函数在 http 请求返回之前看到了 observables。
如何串行执行 HTTP 请求以使输出如下所示?
提交请求 http响应 提交请求 http响应 提交请求 http响应 提交请求 http响应 提交请求 http响应 完全的【问题讨论】:
-
附带说明,您可以通过合并运算符来降低复杂性。
map+flatMap=>flatMap,map+concatAll=>concatMap. -
谢谢,我已经更新了示例以反映这一点
标签: javascript reactive-programming rxjs