You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Reactive Extensions中IGroupedObservable采样结果与预期不符问题

问题根因

你碰到的行为差异是GroupBy + 分组内Sample的组合和全局Sample有两个核心逻辑不一致导致的:

  1. 采样定时器启动时机错位
    全局直接调用Sample(TimeSpan.FromTicks(50))时,测试调度器启动就完成了对Sample的订阅,定时器会从0时刻开始按50tick间隔全局对齐采样,采样点固定为50、100、150、200、250、300、350、400。
    而GroupBy要等到第一个元素到达(230tick)时才会创建对应分组的Observable,此时SelectMany才会订阅分组流、启动分组内的Sample定时器,采样点就变成了从230tick开始每50tick触发,第一个采样点就是280tick,和全局对齐的250tick错位。
  2. 完成信号传播逻辑不同
    全局Sample收到源的OnCompleted信号时,会等待下一个采样点到来时,再推送最后一个缓存值并完成,所以能拿到400tick的输出。
    但GroupBy收到源的OnCompleted信号时,会立刻给所有分组流发送OnCompleted,分组内的Sample收到完成信号后会立刻推送当前缓存的最新值、直接结束,不会等待下一个采样点,所以拿到的是380tick的完成和输出,不会等到400tick。

解决方案

要让分组后的采样行为和全局Sample完全一致,可以改用Sample的采样触发器重载,提前创建全局对齐的采样信号,所有分组共用同一个采样触发流即可:

public class ObservableTests : ReactiveTest
{
    [Fact]
    public void sample_test()
    {
        // Arrange
        var productPrice = new ProductPrice("1", 10);
        var productPrice2 = new ProductPrice("1", 20);
        var productPrice3 = new ProductPrice("1", 30);
        var productPrice4 = new ProductPrice("1", 40);
        var productPrice5 = new ProductPrice("1", 50);
        var productPrice6 = new ProductPrice("1", 60);

        var testScheduler = new TestScheduler();
        var observable = testScheduler.CreateHotObservable(
            OnNext(230, productPrice),
            OnNext(260, productPrice2),
            OnNext(280, productPrice3),
            OnNext(340, productPrice4),
            OnNext(360, productPrice5),
            OnNext(380, productPrice6),
            OnCompleted<ProductPrice>(380));

        // 提前创建全局对齐的50tick采样信号
        var sampleTrigger = Observable.Interval(TimeSpan.FromTicks(50), testScheduler);

        // Act
        var result = testScheduler.Start(
            () =>
                observable
                    .GroupBy(value => value.Identifier)
                    .SelectMany(groupedObservable => groupedObservable.Sample(sampleTrigger)));

        result.Messages.AssertEqual(
            OnNext(250, productPrice),
            OnNext(300, productPrice3),
            OnNext(350, productPrice4),
            OnNext(400, productPrice6),
            OnCompleted<ProductPrice>(400));
    }

    private class ProductPrice
    {
        public ProductPrice(string identifier, decimal price)
        {
            this.Identifier = identifier;
            this.Price = price;
        }

        public string Identifier { get; }

        public decimal Price { get; }
    }
}

内容的提问来源于stack exchange,提问作者Vincent Rutten

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 11:45:05