【问题标题】:rx.net merge reset streamrx.net 合并重置流
【发布时间】:2020-01-29 14:00:57
【问题描述】:

我们有一个正在运行的服务,它处理从系统 x 到系统 y 的消息。 它基本上看起来如下:

aSystem.Messages.Subscribe(message => 
{
   try
   {
      ProcessMessage(message);
   }
   catch(Exception ex)
   {
      _logger.LogFatal(ex.Message);
   }
})

问题是我们每秒至少收到一条消息,并且我们的LogFatal 消息被配置为发送电子邮件。结果,邮箱在某个瞬间爆炸了。

通过添加将保存最后一个时间戳的自定义 Logging 类来“改进”代码。根据该时间戳,消息是否会记录。

这看起来很麻烦,我认为这是Rx.NET 的完美场景。我们需要的是:

  • 1) 如果字符串发生变化,记录日志
  • 2) 超过一定时间后记录

我尝试的是以下内容:

var logSubject = new Subject<string>();
var logMessagesChangedStream = logSubject.DistinctUntilChanged(); // log on every message change
var logMessagesSampleStream = logSubject.Sample(TimeSpan.FromSeconds(10)); // log at least every 10 seconds

var subscription = logMessagesChangedStream.Merge(logMessagesSampleStream).Subscribe(result =>
{
    _logger.LogFatal(result);
});

aSystem.Messages.Subscribe(message => 
{
   try
   {
      ProcessMessage(message);
   }
   catch(Exception ex)
   {
      logSubject.OnNext(ex.Message);
   }
})

看起来它正在工作,但这将记录消息两次,一次用于DistinctUntilChanged,一次用于Sample。因此,如果其中一个流发出了值,我应该以某种方式重置流。他们独立工作完美,一旦合并,他们应该互相倾听;-)

【问题讨论】:

  • 想解释一下否决票?有什么不清楚的地方?
  • 不知道为什么投反对票。支持支持。
  • 您想每十秒(重复)记录最后一个值吗?
  • @Asti 否,仅当有新值时。但最大。每十秒一次。

标签: c# system.reactive rx.net


【解决方案1】:

有一个模棱两可的运算符 Amb 竞争两​​个序列以查看哪个先获胜。

Observable.Amb(logMessagesChangedStream, logMessagesSampleStream) 

获胜的流继续传播到最后 - 我们不希望这样。我们有兴趣再次开始竞争下一个价值。让我们这样做:

    Observable.Amb(logMessagesChangedStream, logMessagesSampleStream)
      .Take(1)
      .Repeat()

现在最后一个问题是DistinctUntilChanged 每次我们重新开始比赛时都会失去它的状态,它的行为是立即推送它获得的第一个值。让我们把它变成一个 hot observable 来解决这个问题。

logSubject.DistinctUntilChanged().Publish();

把它们放在一起:

var logSubject = new Subject<string>();
var logMessagesChangedStream = logSubject.DistinctUntilChanged().Publish(); // log on every message change
var logMessagesSampleStream = logSubject.Sample(TimeSpan.FromSeconds(5)); // log at least every 5 seconds    

var oneof =
        Observable
          .Amb(logMessagesChangedStream, logMessagesSampleStream)
          .Take(1)
          .Repeat(); 

   logMessagesChangedStream.Connect();
   oneof.Timestamp().Subscribe(c => Console.WriteLine(c));

将 10 秒更改为 5,because

测试

    Action<string> onNext = logSubject.OnNext;

    onNext("a");
    onNext("b");
    Delay(1000, () => { onNext("c"); onNext("c"); onNext("c"); onNext("d"); });
    Delay(3000, () => { onNext("d"); onNext("e"); });
    Delay(6000, () => { onNext("e"); });
    Delay(10000, () => { onNext("e"); });

输出

a@1/30/2020 12:17:52 PM +00:00
b@1/30/2020 12:17:52 PM +00:00
c@1/30/2020 12:17:53 PM +00:00
d@1/30/2020 12:17:53 PM +00:00
e@1/30/2020 12:17:55 PM +00:00
e@1/30/2020 12:18:00 PM +00:00
e@1/30/2020 12:18:05 PM +00:00

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-03
    • 1970-01-01
    • 1970-01-01
    • 2020-01-01
    相关资源
    最近更新 更多