【问题标题】:Transform Observable if other Observables emmited mapping function如果其他 Observable 发出映射函数,则转换 Observable
【发布时间】:2015-04-26 18:24:35
【问题描述】:

我正在创建一个游戏,其中有一个可观察的事件流 X 代表制造商交付的产品。还有一些外部事件(我们称之为变形金刚)以各种方式和不同的时间段影响制造的性能。我想通过其他可观察对象来表示这一点,这些可观察对象发出一个转换 X 的函数,并且应该将其应用于每个 X,直到 Transformer 的 OnComplete 为止。变形金刚的数量是未知的 - 它们是由用户操作(如设备购买)或随机生成(如设备故障)创建的。

我想我需要一个 IObservable<IObservable<Func<X,X>>>JoinZip,还有别的?)和 IObservable<X> 来执行此操作。你能帮我解决这个问题吗? Observable.CombineLatest 几乎是我需要的,但它需要 IEnumerable<IObservable<T>>

如果我的描述不清楚,这里有一个大理石图:

用更抽象的术语来说,我需要的是非常类似于矩阵的转置,但不是List<List<T>>,而是IObservable<IObservable<T>>

【问题讨论】:

  • 您使用哪种语言?
  • 目前是 C#,但我很可能在不久的将来将其移植到 Java。
  • 从您的大理石图中看起来,转换不会影响事件的时间。准确吗?
  • 没错——这些只是映射函数,就像普通的 Select() 一样。
  • 申请顺序重要吗?换句话说,您的变换函数是否相互交换?即给定两个变换 f 和 g 以及一个事件 x,f(g(x)) = g(f(x)) 吗?如果它很重要,您将如何订购转换?指定吗?最后申请,最后申请,最后申请等。

标签: system.reactive


【解决方案1】:

假设您的转换器在 int 上工作并且您的 observable 命名如下:

IObservable<IObservable<Func<int, int>>> transformerObservables = null;
IObservable<int> values = null;

我会先将 Transformers 的 Observable 转化为 Transformers 数组的 Observable,也就是

IObservable<IObservable<Func<int, int>>> -> IObservable<<Func<int, int>>[]>

首先,我们最终希望在列表中添加和删除函数,并确保删除正确的转换器,我们必须覆盖 Func<...> 上的常用比较机制。所以我们...

var transformerArrayObservable = transformerObservables
    // ...attach each transformer the index of the observable it came from:        
    .Select((transformerObservable, index) => transformerObservable
        .Select(transformer => Tuple.Create(index, transformer))
        // Then, materialize the transformer sequence so we get noticed when the sequence terminates.
        .Materialize()
        // Now the fun part: Make a scan, resulting in an observable of tuples
        // that have the previous and current transformer
        .Scan(new
        {
            Previous = (Tuple<int, Func<int, int>>)null,
            Current = (Tuple<int, Func<int, int>>)null
        },
        (tuple, currentTransformer) => new
        {
            Previous = tuple.Current,
            Current = currentTransformer.HasValue
                ? currentTransformer.Value
                : (Tuple<int, Func<int, int>>)null
        }))
        // Merge these and do another scan, this time adding and removing
        // the transformers from a list.
        .Merge()
        .Scan(
            new Tuple<int, Func<int, int>>[0],
            (array, tuple) =>
            {
                //Expensive! Consider taking a dependency on immutable collections here!
                var list = array.ToList();

                if (tuple.Previous != null)
                    list.Remove(tuple.Previous);

                if (tuple.Current != null)
                    list.Add(tuple.Current);

                return list.ToArray();
            })
            // Extract only the actual functions
        .Select(x => x.Select(y => y.Item2).ToArray())
        // Finally, to make sure that values are passed even when no transformer has been observed
        // start this sequence with the neutral transformation.
        // IMPORTANT: You should test what happens when the first value is oberserved very quickly. There might be timing issues.
        .StartWith(Scheduler.Immediate, new[] { new Func<int, int>[0]});

现在,您将需要一个 Rx 中不可用的运算符,称为 CombineVeryLatest。看看here

var transformedValues = values
    .CombineVeryLatest(transformerArrayObservable, (value, transformers) =>
    {
        return transformers
            .Aggregate(value, (current, transformer) => transformer(current));
    });

你应该完成了。我敢肯定,有一些性能要获得,但你会明白的。

【讨论】:

    【解决方案2】:

    哇,那真是令人费解,但我认为我有一些可行的方法。首先,我创建了一个扩展方法来将IObservable&lt;IObservable&lt;Func&lt;T, T&gt;&gt; 转换为IObservable&lt;IEnumerable&lt;Func&lt;T, T&gt;&gt;。扩展方法的运行假设每个 observable 在完成之前只会产生一个 Func&lt;T, T&gt;

    public static class MoreReactiveExtensions
    {
        public static IObservable<IEnumerable<Func<T, T>>> ToTransformations<T>(this IObservable<IObservable<Func<T, T>>> source)
        {
            return
                Observable
                // Yield an empty enumerable first.
                .Repeat(Enumerable.Empty<Func<T, T>>(), 1)
                // Then yield an updated enumerable every time one of 
                // the transformation observables yields a value or completes.
                .Concat(                                    
                    source
                    .SelectMany((x, i) => 
                        x
                        .Materialize()
                        .Select(y => new 
                            { 
                                Id = i, 
                                Notification = y 
                            }))
                    .Scan(
                        new List<Tuple<int, Func<T, T>>>(),
                        (acc, x) => 
                        {
                            switch(x.Notification.Kind)
                            {
                                // If an observable compeleted then remove
                                // its corresponding function from the accumulator.
                                case NotificationKind.OnCompleted:
                                    acc = 
                                        acc
                                        .Where(y => y.Item1 != x.Id)
                                        .ToList();
                                    break;
                                // If an observable yield a new Func then add
                                // it to the accumulator.
                                case NotificationKind.OnNext:
                                    acc = new List<Tuple<int, Func<T, T>>>(acc) 
                                        { 
                                            Tuple.Create(x.Id, x.Notification.Value) 
                                        };
                                    break;
                                // Do something with exceptions here.
                                default:
                                    // Do something here
                                    break;
                            }
                            return acc;
                        })
                    // Select an IEnumerable<Func<T, T>> here.
                    .Select(x => x.Select(y => y.Item2)));
        }
    }
    

    然后,给定以下变量:

    IObservable<IObservable<Func<int, int>>> transformationObservables
    IObservable<int> products`
    

    我是这样使用它的:

    var transformations =
        transformationObservables
        .ToTransformations()
        .Publish()
        .RefCount();
    
    IObservable<int> transformedProducts=
        transformations
        .Join(
            products,
            t => transformations,
            i => Observable.Empty<int>(),
            (t, i) => t.Aggregate(i, (ii, tt) => tt.Invoke(ii)))
    

    根据我的测试,结果似乎是正确的。

    【讨论】:

      【解决方案3】:

      受到this answer 的启发,我最终得到了这个:

              Output = Input
                  .WithLatestFrom(
                      transformations.Transpose(),
                      (e, fs) => fs.Aggregate(e, (x, f) => f(x)))
                  .SelectMany(x => x)
                  .Publish();
      

      Transpose 和 WithLatestFrom 运算符定义为:

          public static IObservable<IObservable<T>> Transpose<T>(this IObservable<IObservable<T>> source)
          {
              return Observable.Create<IObservable<T>>(o =>
              {
                  var latestValues = new Dictionary<IObservable<T>, T>();
                  var result = new BehaviorSubject<IObservable<T>>(Observable.Empty<T>());
      
                  source.Subscribe(observable =>
                  {
                      observable.Subscribe(t =>
                      {
                          latestValues[observable] = t;
                          result.OnNext(latestValues.ToObservable().Select(kv => kv.Value));
                      }, () =>
                      {
                          latestValues.Remove(observable);
                      });
                  });
      
                  return result.Subscribe(o);
              });
          }
      
          public static IObservable<R> WithLatestFrom<T, U, R>(
              this IObservable<T> source,
              IObservable<U> other,
              Func<T, U, R> combine)
          {
              return Observable.Create<R>(o =>
              {
                  var current = new BehaviorSubject<U>(default(U));
                  other.Subscribe(current);
                  return source.Select(s => combine(s, current.Value)).Subscribe(o);
              });
          }
      

      这是检查行为的单元测试:

          [TestMethod]
          public void WithLatestFrom_ShouldNotDuplicateEvents()
          {
              var events = new Subject<int>();
      
              var add1 = new Subject<Func<int, int>>();
              var add2 = new Subject<Func<int, int>>();
              var transforms = new Subject<IObservable<Func<int, int>>>();
      
              var results = new List<int>();
      
              events.WithLatestFrom(
                      transforms.Transpose(),
                      (e, fs) => fs.Aggregate(e, (x, f) => f(x)))
                  .SelectMany(x => x)
                  .Subscribe(results.Add);
      
      
              events.OnNext(1);
              transforms.OnNext(add1);
              add1.OnNext(x => x + 1);
              events.OnNext(1); // 1+1 = 2
              transforms.OnNext(add2);
              add2.OnNext(x => x + 2);
              events.OnNext(1); // 1+1+2 = 4
              add1.OnCompleted();
              events.OnNext(1); // 1+2 = 3
              add2.OnCompleted();
              events.OnNext(1);
      
              CollectionAssert.AreEqual(new int[] { 1, 2, 4, 3, 1 }, results);
          }
      

      【讨论】:

      • 我认为,如果您使用CombineLatestOutput 将在任一可观察对象产生值时产生值。意思是,如果transformations 产生一个介于Input 之间的值,那么Output 将产生一个值。根据您的弹珠图,这似乎不是所需的行为。
      • 我认为我在CombineLatest 中看到的另一个问题是OutputInputtransformations 都产生了它们的第一个值之前不会开始产生值。因此,如果您有来自Input 的“产品”流,但transformations 没有产生任何东西,那么Output 将不会产生任何东西。
      • 感谢您注意到这一点。我将通过修复更新答案。
      【解决方案4】:

      将 Transformer 表示为流有意义吗?

      既然添加一个新的Transformer只能改造未来的Events,为什么不维护一些活跃的Transformer集合,然后当一个新的Event进来的时候,你可以应用所有当前的Transformer呢?

      当 Transformer 不再处于活动状态时,它会从集合中移除或标记为非活动状态。

      【讨论】:

        猜你喜欢
        • 2021-09-05
        • 1970-01-01
        • 1970-01-01
        • 2019-12-18
        • 2021-12-03
        • 2023-03-13
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多