【问题标题】:Rx.js wait for callback to completeRxjs 等待回调完成
【发布时间】: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


【解决方案1】:

这里的问题可能是当你使用rx.Observable.fromCallback时,你传入参数的函数会立即执行。返回的 observable 将保存稍后传递给回调的值。为了更好地了解正在发生的事情,您应该使用稍微复杂一点的模拟:为您的请求编号,让它们返回您可以通过订阅观察到的实际(每个请求不同)结果。

我的假设发生在这里:

  • take(5) 发出 5 个值
  • map 发出 5 个日志消息,执行 5 个函数并传递 5 个 observables
  • 这 5 个 observables 由concatAll 处理,这些 observables 发出的值将按预期顺序排列。您在这里订购的是函数调用的结果,而不是函数本身的调用。

为了实现您的目标,您需要仅在concatAll 订阅而不是在创建时调用您的可观察工厂 (rx.Observable.fromCallback)。为此,您可以使用deferhttps://github.com/Reactive-Extensions/RxJS/blob/master/doc/api/core/operators/defer.md

所以你的代码会变成:

rxNode.fromReadableStream(fs.createReadStream('./path-to-file'))
  .map(t => t.toString().split('\r\n'))
  .flatMap(t => t)
  .take(5)
  .map(t => {
    console.log('Submitting request');

    return Observable.defer(function(){return doHttpRequest(t);})
  })
  .concatAll()
  .subscribe(results => {
    console.log(results);
  }, err => {
    console.error('Error', err);
  }, () => {
    console.log('Completed');
  });

你可以在这里看到一个类似的问题,并有一个很好的解释:How to start second observable *only* after first is *completely* done in rxjs

您的日志可能仍会显示 5 条连续的“提交请求”消息。但是您的请求应该按照您的意愿一个接一个地执行。

【讨论】:

    猜你喜欢
    • 2016-02-17
    • 2021-11-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-12-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多