【问题标题】:Accumulating and resetting values in a stream在流中累积和重置值
【发布时间】:2016-02-07 15:26:51
【问题描述】:

我正在玩响应式编程,使用 RxJS,但偶然发现了一些我不知道如何解决的问题。

假设我们实现了一台自动售货机。你投入一枚硬币,选择一个物品,机器会分发一个物品并返回找零。我们假设价格始终为 1 美分,因此插入四分之一(25 美分)应返回 24 美分,依此类推。

“棘手”的部分是我希望能够处理用户插入 2 个硬币然后选择一个项目这样的情况。或者选择一个项目而不插入硬币。

将插入的硬币和选定的项目实现为流似乎很自然。然后我们可以在这两个动作之间引入某种依赖关系——合并或压缩或合并最新的。

但是,我很快遇到了一个问题,我希望将硬币累积到分发物品之前,而不是进一步。 AFAIU,这意味着我不能使用sumscan,因为在某个时候无法“重置”之前的累积。

这是一个示例图:

coins: ---25---5-----10------------|->
acc:   ---25---30----40------------|->
items: ------------foo-----bar-----|->
combined: ---------30,foo--40,bar--|->
change:------------29------39------|->

以及对应的代码:

this.getCoinsStream()
  .scan(function(sum, current) { return sum + current })
  .combineLatest(this.getSelectedItemsStream())
  .subscribe(function(cents, item) {
    dispenseItem(item);
    dispenseChange(cents - 1);
  });

插入 25 和 5 美分,然后选择“foo”项目。累积硬币然后组合最新的将导致“foo”与“30”(这是正确的)组合,然后是“bar”与“40”(不正确;应该是“bar”和“10”)。

我查看了all of the methods 进行分组和过滤,但没有看到任何我可以使用的东西。

我可以使用的另一种解决方案是分别积累硬币。但这会在流之外引入状态,我真的很想避免这种情况:

var centsDeposited = 0;
this.getCoinsStream().subscribe(function(cents) {
  return centsDeposited += cents;
});

this.getSelectedItemsStream().subscribe(function(item) {
  dispenseItem(item);
  dispenseChange(centsDeposited - 1);
  centsDeposited = 0;
});

此外,这不允许使流相互依赖,例如等待硬币被插入,直到选定的操作可以返回一个项目。

我是否缺少已经存在的方法?实现这样的事情的最佳方法是什么 - 累积值直到它们需要与另一个流合并的那一刻,但还要等待第一个流中的至少一个值,然后再将其与第二个流中的值合并?

【问题讨论】:

标签: javascript asynchronous reactive-programming rxjs


【解决方案1】:

您可以使用您的 scan/combineLatest 方法,然后使用 first 结束流,然后使用 repeat 以便它“重新开始”流,但您的观察者不会看到它。

var coinStream = Rx.Observable.merge(
    Rx.Observable.fromEvent($('#add5'), 'click').map(5),
    Rx.Observable.fromEvent($('#add10'), 'click').map(10),
    Rx.Observable.fromEvent($('#add25'), 'click').map(25)
    );
var selectedStream = Rx.Observable.merge(
  Rx.Observable.fromEvent($('#coke'), 'click').map('Coke'),
  Rx.Observable.fromEvent($('#sprite'), 'click').map('sprite')
);

var $selection = $('#selection');
var $change = $('#change');

function dispense(selection) {
 $selection.text('Dispensed: ' + selection); 
 console.log("Dispensing Drink: " + selection);   
}

function dispenseChange(change) {
 $change.text('Dispensed change: ' + change); 
 console.log("Dispensing Change: " + change);   
}


var dispenser = coinStream.scan(function(acc, delta) { return acc + delta; }, 0)
                              .combineLatest(selectedStream, 
                                 function(coins, selection) {
                                   return {coins : coins, selection : selection};
                                 })

                              //Combine latest won't emit until both Observables have a value
                              //so you can safely get the first which will be the point that
                              //both Observables have emitted.
                              .first()

                              //First will complete the stream above so use repeat
                              //to resubscribe to the stream transparently
                              //You could also do this conditionally with while or doWhile
                              .repeat()

                              //If you only will subscribe once, then you won't need this but 
                              //here I am showing how to do it with two subscribers
                              .publish();

    
    //Dole out the change
    dispenser.pluck('coins')
             .map(function(c) { return c - 1;})
             .subscribe(dispenseChange);
    
    //Get the selection for dispensation
    dispenser.pluck('selection').subscribe(dispense);
    
    //Wire it up
    dispenser.connect();
<script src="https://ajax.googleapis.com/ajax/libs/jquery/2.1.1/jquery.min.js"></script>
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/4.0.6/rx.all.js"></script>
<button id="coke">Coke</button>
<button id="sprite">Sprite</button>
<button id="add5">5</button>
<button id="add10">10</button>
<button id="add25">25</button>
<div id="change"></div>
<div id="selection"></div>

【讨论】:

  • 啊哈,所以诀窍是用first()“结束”流,然后用repeat()“恢复”?有趣的。但也有点混乱。我希望repeat() 无限重复最后一个值(这是找到的第一个组合)。另一方面,如果重复重新启动整个流,那么它是有道理的。非常酷的方法,一点也不明显:) 谢谢。
  • 正确。一种思考方式是,每次购买苏打水时,您实际上都在开始新的交易(其中将提供新硬币并且在购买苏打水之前不会完成)。您实际上只对用户推送选择支付的第一次时间感兴趣。之后,您想“重复”或重新启动新事务(请注意,如果这令人困惑,您还可以查看 whiledoWhile,它们做类似的事情并且可能更熟悉)
【解决方案2】:

一般来说,您有以下一组方程:

inserted_coins :: independent source
items :: independent source
accumulated_coins :: sum(inserted_coins)
accumulated_paid :: sum(price(items))
change :: accumulated_coins - accumulated_paid
coins_in_machine :: when items : 0, when inserted_coins : sum(inserted_coins) starting after last emission of item

困难的部分是coins_in_machine。您需要根据来自两个来源的一些排放来切换可观察的来源。

function emits ( who ) {
  return function ( x ) { console.log([who, ": "].join(" ") + x);};
}

function sum ( a, b ) {return a + b;}

var inserted_coins = Rx.Observable.fromEvent(document.getElementById("insert"), 'click').map(function ( x ) {return 15;});
var items = Rx.Observable.fromEvent(document.getElementById("item"), 'click').map(function ( x ) {return "snickers";});

console.log("running");

var accumulated_coins = inserted_coins.scan(sum);

var coins_in_machine =
    Rx.Observable.merge(
        items.tap(emits("items")).map(function ( x ) {return {value : x, flag : 1};}),
        inserted_coins.tap(emits("coins inserted ")).map(function ( x ) {return {value : x, flag : 0};}))
        .distinctUntilChanged(function(x){return x.flag;})
        .flatMapLatest(function ( x ) {
                   switch (x.flag) {
                     case 1 :
                       return Rx.Observable.just(0);
                     case 0 :
                       return inserted_coins.scan(sum, x.value).startWith(x.value);
                   }
                 }
    ).startWith(0);

coins_in_machine.subscribe(emits("coins in machine"));

jsbin : http://jsbin.com/mejoneteyo/edit?html,js,console,output

[更新]

解释:

  • 我们将 insert_coins 流与 items 流合并,同时为它们附加一个标志,以了解当我们在合并流中接收到值时发出两者中的哪一个
  • 当它是项目流发射时,我们想把0放在coins_in_machine中。当它是insert_coins 时,我们想要对传入的值求和,因为该总和将代表机器中新的硬币数量。这意味着insert_coins 的定义在之前定义的逻辑下从一个流切换到另一个流。该逻辑是在switchMapLatest 中实现的。
  • 我使用switchMapLatest 而不是switchMap,否则coins_in_machine 流将继续接收来自以前切换的流的发射,即重复发射,因为最终只有两个流进出我们切换.如果可以的话,我会说这是我们需要的关闭和切换。
  • switchMapLatest 必须返回一个流,所以我们跳过箍来制作一个发出 0 和 never 结束的流(并且不会阻塞计算机,因为在这种情况下使用 repeat 运算符会)
  • 我们跳过了一些额外的环节以使inserted_coins 发出我们想要的值。我的第一个实现是inserted_coins.scan(sum,0),但从未奏效。关键是我发现这很棘手,当我们到达流程中的那个点时,inserted_coins 已经发出了作为总和一部分的值之一。该值是作为flatMapLatest 的参数传递的值,但它不再在源中,因此在事实之后调用scan 不会得到它,因此有必要从flatMapLatest 获取该值和重构正确的行为。

【讨论】:

  • 很好的问题分解和有趣的解决方案,谢谢!因此,使用distinctUntilChanged 并通过标志区分流是一种选择。但是该死的......这真的很难理解:/
  • 添加了一些解释。是的,这不是微不足道的,我花了几个小时才弄清楚我最初的幼稚实现中出了什么问题,但是在处理热的高阶 observables 时要小心发射时间是一个很好的见解。
  • 不确定我理解为什么在这种情况下你需要never().startWith()just 应该做同样的事情,因为下游永远不会“看到”中间流的完成。
  • 刚刚进行了更改并进行了测试。你说得对。已更正。
【解决方案3】:

您还可以使用Window 将多个硬币事件组合在一起,并使用项目选择作为窗口边界。

接下来我们可以使用 zip 来获取 item 的值。

请注意,我们会立即尝试分发物品。所以用户在决定一个项目之前确实必须插入硬币。

请注意,出于安全原因,我决定同时发布 selectedStreamdispenser,我们不希望在我们构建查询时引发事件并且 zip 变得不平衡的竞争条件。这将是一种非常罕见的情况,但请注意,当我们的源是冷 Observable 时,它​​们几乎会在我们订阅后立即开始生成,我们必须使用 Publish 来保护自己。

(无耻窃取 paulpdaniels 示例代码)。

var coinStream = Rx.Observable.merge(
    Rx.Observable.fromEvent($('#add5'), 'click').map(5),
    Rx.Observable.fromEvent($('#add10'), 'click').map(10),
    Rx.Observable.fromEvent($('#add25'), 'click').map(25)
);

var selectedStream = Rx.Observable.merge(
    Rx.Observable.fromEvent($('#coke'), 'click').map('Coke'),
    Rx.Observable.fromEvent($('#sprite'), 'click').map('Sprite')
).publish();

var $selection = $('#selection');
var $change = $('#change');

function dispense(selection) {
    $selection.text('Dispensed: ' + selection); 
    console.log("Dispensing Drink: " + selection);   
}

function dispenseChange(change) {
    $change.text('Dispensed change: ' + change); 
    console.log("Dispensing Change: " + change);   
}

// Build the query.
var dispenser = Rx.Observable.zip(
    coinStream
        .window(selectedStream)
        .flatMap(ob => ob.reduce((acc, cur) => acc + cur, 0)),
    selectedStream,
    (coins, selection) => ({coins : coins, selection: selection})
).filter(pay => pay.coins != 0) // Do not give out items if there are no coins.
.publish();

var dispose = new Rx.CompositeDisposable(
    //Dole out the change
    dispenser
        .pluck('coins')
        .map(function(c) { return c - 1;})
        .subscribe(dispenseChange),

    //Get the selection for dispensation
    dispenser
        .pluck('selection')
        .subscribe(dispense),

    //Wire it up
    dispenser.connect(),
    selectedStream.connect()
);
<script src="https://ajax.googleapis.com/ajax/libs/jquery/2.1.1/jquery.min.js"></script>
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/4.0.6/rx.all.js"></script>
<button id="coke">Coke</button>
<button id="sprite">Sprite</button>
<button id="add5">5</button>
<button id="add10">10</button>
<button id="add25">25</button>
<div id="change"></div>
<div id="selection"></div>

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-27
    • 2020-11-18
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多