【问题标题】:How to buffer observables in pairs and execute them by the pair?如何成对缓冲 observables 并成对执行它们?
【发布时间】:2017-10-18 19:26:56
【问题描述】:

我有一个 xhr 请求,它正在获取一个数组,我使用该数组执行后续 xhr 请求,如下所示:

const Rx = require('rxjs/Rx');
const fetch = require('node-fetch');

const url = `url`;

// Get array of tables
const tables$ = Rx.Observable
  .from(fetch(url).then((r) => r.json()));

// Get array of columns
const columns$ = (table) => {
  return Rx.Observable
    .from(fetch(`${url}/${table.TableName}/columns`).then(r => r.json()));
};

tables$
  .mergeMap(tables => Rx.Observable.forkJoin(...tables.map(columns$)))    
  .subscribe(val => console.log(val));

我想以块的形式执行列请求,这样请求就不会立即发送到服务器。

这个 SO 问题有点相同但不完全:Rxjs: Chunk and delay stream?

现在我正在尝试这样的事情:

tables$
  .mergeMap(tables => Rx.Observable.forkJoin(...tables.map(columns$)))
  .flatMap(e => e)
  .bufferCount(4)
  .executeTheChunksSerial(magic)
  .flatMap(e => e)
  .subscribe(val => console.log(val));

但我无法理解如何连续执行这些块......

【问题讨论】:

  • 看起来你有不止一个高阶 Observable,这使得这很难理解。也许你可以使用.concatMap(buffered => Observable.forkJoin(buffered)) 而不是.executeTheChunksSerial(magic),但我不确定。

标签: rxjs rxjs5


【解决方案1】:

您可以利用mergeMapconcurrency 参数来同时向您的服务器获取最大x 个请求:

const getTables = Promise.resolve([{ tableName: 'foo' },{ tableName: 'bar' },{ tableName: 'baz' }]);
const getColumns = (table) => Rx.Observable.of('a,b,c')
  .do(_ => console.log('getting columns for table: ' + table))
  .delay(250);
      
Rx.Observable.from(getTables)
  .mergeAll()
  .mergeMap(
    table => getColumns(table.tableName),
    (table, columns) => ({ table, columns }),
    2)
  .subscribe(console.log)
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/5.4.3/Rx.js"></script>

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-02-26
    • 2013-07-27
    • 2013-09-21
    • 1970-01-01
    • 1970-01-01
    • 2016-01-08
    相关资源
    最近更新 更多