【问题标题】:Take first elements of stream after previous element matches condition在前一个元素匹配条件之后获取流的第一个元素
【发布时间】:2018-08-21 08:44:10
【问题描述】:

我是响应式扩展 (rx) 的新手,并尝试在 .NET 中执行以下操作(但在 JS 和其他语言中应该相同。):

我有一个带有包含字符串和布尔属性的对象的流传入。流将是无限的。我有以下条件:

  • 应始终打印第一个对象。
  • 现在应该跳过所有传入的对象,直到有布尔属性设置为“true”的对象到达。
  • 当一个对象到达时 bool 属性设置为“true”,应该跳过这个对象,但应该打印下一个对象(不管属性是什么)。
  • 现在这样下去,应该打印属性设置为 true 的对象之后的每个对象。

例子:

("one", false)--("two", true)--("three", false)--("four", false)--("five", true)--("six", true)--("seven", true)--("eight", false)--("nine", true)--("ten", false)

预期结果:

"one"--"three"--"six"--"seven"--"eight"--"ten"

注意“六”和“七”已被打印,因为它们跟随一个属性设置为 true 的对象,即使它们自己的属性也设置为“true”。

简单的 .NET 程序来测试它:

using System;
using System.Reactive.Linq;
using System.Threading;

namespace ConsoleApp1
{
    class Program
    {
        class Result
        {
            public bool Flag { get; set; }
            public string Text { get; set; }
        }

        static void Main(string[] args)
        {               
            var source =
               Observable.Create<Result>(f =>
               {
                   f.OnNext(new Result() { Text = "one", Flag = false });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "two", Flag = true });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "three", Flag = false });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "four", Flag = false });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "five", Flag = true });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "six", Flag = true });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "seven", Flag = true });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "eight", Flag = false });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "nine", Flag = true });
                   Thread.Sleep(1000);
                   f.OnNext(new Result() { Text = "ten", Flag = false });

                   return () => Console.WriteLine("Observer has unsubscribed");
               });
        }
    }
}

我尝试使用 .Scan 和 .Buffer 扩展,但我不知道如何在我的场景中使用它们。

性能当然要尽可能好,因为最终流是无限的。

【问题讨论】:

    标签: c# .net rxjs reactive-programming system.reactive


    【解决方案1】:

    试试这个方法:

    var results = new[]
    {
        new Result() { Text = "one", Flag = false },
        new Result() { Text = "two", Flag = true },
        new Result() { Text = "three", Flag = false },
        new Result() { Text = "four", Flag = false },
        new Result() { Text = "five", Flag = true },
        new Result() { Text = "six", Flag = true },
        new Result() { Text = "seven", Flag = true },
        new Result() { Text = "eight", Flag = false },
        new Result() { Text = "nine", Flag = true },
        new Result() { Text = "ten", Flag = false },
    };
    
    IObservable<Result> source =
        Observable
            .Generate(
                0, x => x < results.Length, x => x + 1,
                x => results[x],
                x => TimeSpan.FromSeconds(1.0));
    

    以上只是以比您的Observable.Create&lt;Result&gt; 方法更惯用的方式生成source。

    现在是query:

    IObservable<Result> query =
        source
            .StartWith(new Result() { Flag = true })
            .Publish(ss =>
                ss
                    .Skip(1)
                    .Zip(ss, (s1, s0) =>
                        s0.Flag
                        ? Observable.Return(s1) 
                        : Observable.Empty<Result>())
                    .Merge());
    

    这里使用.Publish 允许源observable 只有一个订阅,但可以在.Publish 方法中多次使用它。然后可以使用标准的Skip(1).Zip 方法来检查正在生成的后续值。

    这是输出:


    从 Shlomo 获得灵感,这是我使用 .Buffer(2, 1) 的方法:

    IObservable<Result> query2 =
        source
            .StartWith(new Result() { Flag = true })
            .Buffer(2, 1)
            .Where(rs => rs.First().Flag)
            .SelectMany(rs => rs.Skip(1));
    

    【讨论】:

    • 工作就像一个魅力,谢谢!您能否进一步解释一下您为什么使用发布,因为它让我有点困惑。 :)
    • @TobiasvonFalkenhayn - 你不确定我的回答中对.Publish 的使用的描述有什么特别之处吗?
    • 将 Console.Writeline("Nasty Side-Effect."); 添加到源 observable 的开头。使用Publish,您将看到该消息一次。没有它,你会看到它两次。
    • 每次你说source. 时,都会创建一个新订阅,它会从头开始再次执行整个 observable。 Publish 就像一个路由器,它“复制”它收到的通知并传递多个副本。
    • @TobiasvonFalkenhayn - 你有点偏离轨道。当您使用新运算符时,不会创建新订阅。当您使用source 时。通常,当您使用对source 的单独引用时,您有单独的订阅。使用x.Publish(y =&gt; ...) 会更改规则,这样您就可以在x 上仅订阅一次,您可以根据需要多次使用y。
    【解决方案2】:

    这里有很多方法可以做到:

    var result1 = source.Publish(_source => _source
        .Zip(_source.Skip(1), (older, newer) => (older, newer))
        .Where(t => t.older.Flag == true)
        .Select(t => t.newer)
        .Merge(_source.Take(1))
        .Select(r => r.Text)
    );
    
    var result2 = source.Publish(_source => _source
        .Buffer(2, 1)
        .Where(l => l[0].Flag == true)
        .Select(l => l[1])
        .Merge(_source.Take(1))
        .Select(l => l.Text)
    );
    
    var result3 = source.Publish(_source => _source
        .Window(2, 1)
        .SelectMany(w => w
            .TakeWhile((r, index) => (index == 0 && r.Flag) || index == 1)
            .Skip(1)
        )
        .Merge(_source.Take(1))
        .Select(l => l.Text)
    );
    
    var result4 = source
        .Scan((result: new Result {Flag = true, Text = null}, emit: false), (state, r) => (r, state.result.Flag))
        .Where(t => t.emit)
        .Select(t => t.result.Text);
    

    我偏爱扫描版,但实际上,这取决于你。

    【讨论】:

    • 我自己喜欢.Buffer(2, 1)。我会删除.Merge(_source.Take(1)) 并亲自添加.StartWith(new Result() { Flag = true }),但是嘿,有很多选择。
    • 感谢您的帮助。 :) 你能给我一个链接或为什么你使用“发布”吗?我不明白这与选择器功能的工作原理是什么。
    【解决方案3】:

    感谢@Picci 的回答,我找到了一种方法:

    Func<bool, Action<Result>> printItem = print =>
                    {
                        return data => {
                            if(print)
                            {
                                Console.WriteLine(data.Text);
                            }
                        };
                    };
    
    var printItemFunction = printItem(true);
    
    source.Do(item => printItemFunction(item))
          .Do(item => printItemFunction = printItem(item.Flag))
          .Subscribe();
    

    但是,我不太确定这是否是最好的方法,因为不使用 Subscribe() 但副作用对我来说似乎有点奇怪。最后,我不仅想打印值,还想用它调用 Web 服务。

    【讨论】:

      【解决方案4】:

      这就是我在 TypeScript 中编码的方式

      const printItem = (print: boolean) => {
          return (data) => {
              if (print) {
                  console.log(data);
              }
          };
      }
      
      let printItemFunction = printItem(true);
      
      from(data)
      .pipe(
          tap(item => printItemFunction(item.data)),
          tap(item => printItemFunction = printItem(item.printBool))
      )
      .subscribe()
      

      基本思想是使用更高级别的函数printItem,它返回一个知道是否打印以及打印什么的函数。返回的函数存储在变量printItemFunction中。

      对于源 Observable 发出的每一项,首先要做的就是执行printItemFunction 传递源 Observable 通知的数据。

      第二件事是评估printItem函数并将结果存储在变量printItemFunction中,以便为后续通知做好准备。

      在程序开始时,printItemFunction 被初始化为true,所以总是打印第一项。

      我不熟悉C#给你一个.NET的答案

      【讨论】:

      • 非常感谢。我设法在.NET 中做同样的事情。见我上面的回答。但是,我不太确定这是否是最好的方法,因为不使用 Subscribe() 但副作用对我来说似乎有点奇怪。最后,我不仅想打印值,还想用它调用 Web 服务。
      • 一旦解决方案的关键元素明确,您就可以混合副作用并订阅,这也是定义副作用的另一种方式,随心所欲。我在代码中使用它们以试图明确关键步骤。
      • @Picci - 我建议尽可能避免副作用。
      • 我猜在控制台上打印是个副作用,不是吗?
      • @Picci - 您应该尝试拥有两个当前订阅者。这会导致竞争条件。
      猜你喜欢
      • 2014-05-21
      • 1970-01-01
      • 1970-01-01
      • 2018-07-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多