【问题标题】:Use First() and Repeat() without restarting whole stream RxJS使用 First() 和 Repeat() 而不重新启动整个流 RxJS
【发布时间】:2018-02-04 15:57:14
【问题描述】:

我正在使用RxJS 构建一个交易机器人。为此,我必须将来自套接字连接的股票数据转换为每x 秒发出的蜡烛。

我像这样创建了socketObservable

const subscribeObservable = Observable.fromEventPattern(h => bittrex.websockets.subscribe(['USDT-BTC'], h))
const clientCallBackObservable = Observable.fromEventPattern(h => bittrex.websockets.client(h))

const socketObservable = clientCallBackObservable
  .flatMap(() => subscribeObservable)
  .filter(subscribtionData => subscribtionData && subscribtionData.M === 'updateExchangeState')
  .flatMap(exchangeState => Observable.from(exchangeState.A))
  .filter(marketData => marketData.Fills.length > 0)
  .map(marketData => marketData && marketData.Fills)

效果很好 - 当我连接到客户端时,我 flatMap 到订阅连接。

然后我有导致问题的candleObservable

export const candleObservable = (promise, timeFrame = TIME_FRAME) =>
  promise
    .scan((acc, curr) => [...acc, ...curr])
    .skipWhile(exchangeData => dateDifferenceInSeconds(exchangeData) < timeFrame)
    // take first after skipping
    .first()
    // first will complete the stream, so we repeat it
    .repeat()
    // we create candle data from the timeFrame array
    .map(fillsData => createCandle(fillsData))
    // accumulate candles
    .scan((acc, curr) => [...[acc], curr])

我想要实现的是积累数据,直到我有一个完整的蜡烛,可以是x 秒。然后我想用那个发射并重置扫描功能,所以我开始买一根新蜡烛。然后我创建蜡烛并在另一次扫描中累积它。

我的问题是当我打电话给repeat() 我的socketObservable 也被再次调用。我不知道这是否会导致node-bittrex-api 产生任何开销,但我想避免它。

我曾尝试将累积的蜡烛部分放入 flatMap 或类似的东西中,但无法让其中任何一个工作。

你知道我怎样才能避免repeat() 整个流或另一种制作蜡烛的方法,我可以在其中累积,然后在第一次发射后重置累积器?

【问题讨论】:

  • 是否可以简化问题?只有最少的必要运算符,没有特定领域的东西,比如蜡烛。

标签: node.js ecmascript-6 rxjs


【解决方案1】:

根据您的描述,听起来您有一个可观察的对象,您想根据某些条件将其切成某种桶。通常,将一个流减少为具有较少元素(没有过滤)的另一个流被称为“背压”。在您的具体情况下,听起来您感兴趣的背压运算符是bufferbuffer 操作符可以接受一个 observable 作为作为“关闭选择器”的参数,即当你绑定一个缓冲区并启动一个新缓冲区时,这个 observable 中的排放量可以用来调节。

我建议用buffer 调用替换您的scanskipWhilefirstrepeat,传入一个关闭选择器,该选择器将在您的“TIME_FRAME”到期时产生一个值。这应该很容易使用timer(在固定数量的情况下)或驱动流的去抖动版本(如果您想在数据暂停时停止)表示为可观察的。如果您的缓冲区是严格基于时间的,那么甚至还有一个名为 bufferTime 的专门化 buffer 来处理这个问题。因为您最终会得到一个可观察的数组(而不是原始值),所以您可能希望将最终的 scan 替换为常规数组 reduce

如果没有更简单的示例,很难给出具体的代码。我强烈建议您查阅各种背压操作符的示例代码,看看您是否能找到与您想要实现的目标类似的东西。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-13
    • 2020-06-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多