【问题标题】:Creating a weak subscription to an IObservable创建对 IObservable 的弱订阅
【发布时间】:2011-09-06 15:31:44
【问题描述】:

我想要做的是确保如果对我的观察者的唯一引用是可观察的,它会被垃圾收集并停止接收消息。

假设我有一个控件,上面有一个名为 Messages 的列表框,后面的代码如下:

//Short lived display of messages (only while the user's viewing incoming messages)
public partial class MessageDisplay : UserControl
{
    public MessageDisplay()
    {
        InitializeComponent();
        MySource.IncomingMessages.Subscribe(m => Messages.Items.Add(m));
    }
}

哪个连接到这个源:

//Long lived location for message store
static class MySource
{
    public readonly static IObservable<string> IncomingMessages = new ReplaySubject<string>;
}

我不希望消息显示在它不再可见后长时间保留在内存中。理想情况下,我想要一个小扩展,这样我就可以写了:

MySource.IncomingMessages.ToWeakObservable().Subscribe(m => Messages.Items.Add(m));

我也不想依赖 MessageDisplay 是一个用户控件这一事实,因为我稍后会想要使用 MessageDisplayViewModel 进行 MVVM 设置,这将不是用户控件。

【问题讨论】:

  • 当您知道不再需要 Observable 时,您的代码中是否还有位置?在这种情况下,您可以获取从Subscribe 方法返回的IDisposable,以便在需要时将其删除。
  • @seldon 我可以在这个特定的例子中使用它,只要消息窗口关闭,但我想要一个更通用的方法,这样我可以更广泛地使用这个功能并防止其他程序员使用我的图书馆忘记在某处放置一些东西。我在 MVVMLightToolkit 中看到了一些相关的东西,但不适用于 IObservable,而且我并不真正了解它是如何工作的,而且众所周知,这些东西很难做好。
  • 您可能指的是WeakReference 类,它可用于使垃圾收集器收集实例,即使它们被“引用”。然而,在响应式扩展中有很多操作符处理在某些时候处理 observable。如果当你知道你已经完成了 observable 时,你有一些事件等,也许这些就足够了?
  • 我建议您积极尝试采用反模式。正确完成的 MVVM 与 IDispose 模式一起工作得很好。你应该非常正确地实现这一点。

标签: c# garbage-collection system.reactive weak-references


【解决方案1】:

您可以将代理观察者订阅到持有对实际观察者的弱引用的可观察对象,并在实际观察者不再活动时释放订阅:

static IDisposable WeakSubscribe<T>(
    this IObservable<T> observable, IObserver<T> observer)
{
    return new WeakSubscription<T>(observable, observer);
}

class WeakSubscription<T> : IDisposable, IObserver<T>
{
    private readonly WeakReference reference;
    private readonly IDisposable subscription;
    private bool disposed;

    public WeakSubscription(IObservable<T> observable, IObserver<T> observer)
    {
        this.reference = new WeakReference(observer);
        this.subscription = observable.Subscribe(this);
    }

    void IObserver<T>.OnCompleted()
    {
        var observer = (IObserver<T>)this.reference.Target;
        if (observer != null) observer.OnCompleted();
        else this.Dispose();
    }

    void IObserver<T>.OnError(Exception error)
    {
        var observer = (IObserver<T>)this.reference.Target;
        if (observer != null) observer.OnError(error);
        else this.Dispose();
    }

    void IObserver<T>.OnNext(T value)
    {
        var observer = (IObserver<T>)this.reference.Target;
        if (observer != null) observer.OnNext(value);
        else this.Dispose();
    }

    public void Dispose()
    {
        if (!this.disposed)
        {
            this.disposed = true;
            this.subscription.Dispose();
        }
    }
}

【讨论】:

  • +1 这正是我所希望的。一个快速的问题是,它是否仍然适用于 Subscribe(M=>DoSomethingWithM(M)) 的扩展方法,或者它们是否会提前收集垃圾?作为对此的扩展,这是否是您查询中的最后一件事,或者您是否也可以在此之后执行投影/查询等?
  • 您应该在此之前执行投影和过滤。 Subscribe(M=>DoSomethingWithM(M)) 在内部创建一个包装委托 M=>DoSomethingWithM(M) 的 IObserver。您需要保持内部创建的 IObserver 活动,这是不可能的,因为它是内部的。所以,在考虑之后,我实际上不会再推荐我的答案了。寻找不同的方法。
  • 订阅事件时返回的IDisposable是否允许异步处理?通过包含对 IObserver 的引用并明确 Dispose 上的引用,即使线程上下文不允许它做任何其他事情,实现这种行为似乎应该很容易;如果这样做了,IObservable 可以通过定期清除其 Disposed 订阅者的订阅列表来避免内存泄漏。每次添加订阅或发现订阅已被删除时检查一个订阅者就足够了。
  • 不幸的是,即使对 IObservable 施加这样的要求几乎不会很繁重,但我没有看到任何实际施加这种要求的东西,也不知道现有的实现会在多大程度上遵守。
【解决方案2】:

几年后遇到这个线程...只是想指出Samuel Jack's blog 上确定的解决方案,它向 IObservable 添加了一个名为 WeaklySubscribe 的扩展方法。它使用一种在主体和观察者之间添加垫片的方法,该垫片通过 WeakReference 跟踪目标。这类似于其他人针对事件订阅中的强引用问题提供的解决方案,例如this articlethis solution by Paul Stovell。有一段时间使用了基于 Paul 方法的东西,我喜欢 Samuel 针对弱 IObservable 订阅的解决方案。

【讨论】:

    【解决方案3】:

    还有另一个选项使用weak-event-patterns

    基本上System.Windows.WeakEventManager 有你。

    当您的 ViewModel 依赖于带有事件的服务时使用 MVVM,您可以弱订阅这些服务,从而允许您的 ViewModel 与视图一起被收集,而无需事件订阅使其保持活动状态。

    using System;
    using System.Windows;
    
    class LongLivingSubject
    { 
        public event EventHandler<EventArgs> Notifications = delegate { }; 
    }
    
    class ShortLivingObserver
    {
        public ShortLivingObserver(LongLivingSubject subject)
        { 
            WeakEventManager<LongLivingSubject, EventArgs>
                .AddHandler(subject, nameof(subject.Notifications), Subject_Notifications); 
        }
    
        private void Subject_Notifications(object sender, EventArgs e) 
        { 
        }
    }
    

    【讨论】:

      【解决方案4】:

      这是我的实现(退出简单)

      public class WeakObservable<T>: IObservable<T>
      {
          private IObservable<T> _source;
      
          public WeakObservable(IObservable<T> source)
          {
              #region Validation
      
              if (source == null)
                  throw new ArgumentNullException("source");
      
              #endregion Validation
      
              _source = source;
          }
      
          public IDisposable Subscribe(IObserver<T> observer)
          {
              IObservable<T> source = _source;
              if(source == null)
                  return Disposable.Empty;
              var weakObserver = new WaekObserver<T>(observer);
              IDisposable disp = source.Subscribe(weakObserver);
              return disp;
          }
      }
          public class WaekObserver<T>: IObserver<T>
      {
          private WeakReference<IObserver<T>> _target;
      
          public WaekObserver(IObserver<T> target)
          {
              #region Validation
      
              if (target == null)
                  throw new ArgumentNullException("target");
      
              #endregion Validation
      
              _target = new WeakReference<IObserver<T>>(target);
          }
      
          private IObserver<T> Target
          {
              get
              {
                  IObserver<T> target;
                  if(_target.TryGetTarget(out target))
                      return target;
                  return null;
              }
          }
      
          #region IObserver<T> Members
      
          /// <summary>
          /// Notifies the observer that the provider has finished sending push-based notifications.
          /// </summary>
          public void OnCompleted()
          {
              IObserver<T> target = Target;
              if (target == null)
                  return;
      
              target.OnCompleted();
          }
      
          /// <summary>
          /// Notifies the observer that the provider has experienced an error condition.
          /// </summary>
          /// <param name="error">An object that provides additional information about the error.</param>
          public void OnError(Exception error)
          {
              IObserver<T> target = Target;
              if (target == null)
                  return;
      
              target.OnError(error);
          }
      
          /// <summary>
          /// Provides the observer with new data.
          /// </summary>
          /// <param name="value">The current notification information.</param>
          public void OnNext(T value)
          {
              IObserver<T> target = Target;
              if (target == null)
                  return;
      
              target.OnNext(value);
          }
      
          #endregion IObserver<T> Members
      }
          public static class RxExtensions
      {
          public static IObservable<T> ToWeakObservable<T>(this IObservable<T> source)
          {
              return new WeakObservable<T>(source);
          }
      }
              static void Main(string[] args)
          {
              Console.WriteLine("Start");
              var xs = Observable.Interval(TimeSpan.FromSeconds(1));
              Sbscribe(xs);
      
              Thread.Sleep(2020);
              Console.WriteLine("Collect");
              GC.Collect();
              GC.WaitForPendingFinalizers();
              GC.Collect();
              Console.WriteLine("Done");
              Console.ReadKey();
          }
      
          private static void Sbscribe<T>(IObservable<T> source)
          {
              source.ToWeakObservable().Subscribe(v => Console.WriteLine(v));
          }
      

      【讨论】:

        【解决方案5】:

        关键是要认识到您必须同时传入目标和双参数动作。单参数动作永远不会这样做,因为要么你对你的动作使用弱引用(并且动作会被 GC'd),要么你对你的动作使用强引用,这反过来又对目标有强引用,因此目标无法获得 GC。牢记这一点,以下工作:

        using System;
        
        namespace Closures {
          public static class WeakReferenceExtensions {
            /// <summary> returns null if target is not available. Safe to call, even if the reference is null. </summary>
            public static TTarget TryGetTarget<TTarget>(this WeakReference<TTarget> reference) where TTarget : class {
              TTarget r = null;
              if (reference != null) {
                reference.TryGetTarget(out r);
              }
              return r;
            }
          }
          public static class ObservableExtensions {
        
            public static IDisposable WeakSubscribe<T, U>(this IObservable<U> source, T target, Action<T, U> action)
              where T : class {
              var weakRef = new WeakReference<T>(target);
              var r = source.Subscribe(u => {
                var t = weakRef.TryGetTarget();
                if (t != null) {
                  action(t, u);
                }
              });
              return r;
            }
          }
        }
        

        可观察样本:

        using System;
        using System.Reactive.Subjects;
        
        namespace Closures {
          public class Observable {
            public IObservable<int> ObservableProperty => _subject;
            private Subject<int> _subject = new Subject<int>();
            private int n;
            public void Fire() {
              _subject.OnNext(n++);
            }
          }
        }
        

        用法:

        Class SomeClass {
        
         IDisposable disposable;
        
         public void SomeMethod(Observable observeMe) {
           disposable = observeMe.ObservableProperty.WeakSubscribe(this, (wo, n) => wo.Log(n));
         }
        
          public void Log(int n) {
            System.Diagnostics.Debug.WriteLine("log "+n);
          }
        }
        

        【讨论】:

          【解决方案6】:

          下面的代码灵感来自 dtb 的原始帖子。唯一的变化是它作为 IDisposable 的一部分返回对观察者的引用。这意味着只要您保留对您在链末端取出的 IDisposable 的引用,对 IObserver 的引用就会保持活动状态(假设所有一次性用品都保留对一次性用品的引用)。这允许使用诸如Subscribe(M=&gt;DoSomethingWithM(M)) 之类的扩展方法,因为我们保留了对隐式构造的 IObserver 的引用,但我们没有保留从源到 IObserver 的强引用(这会产生内存泄漏)。

          using System.Reactive.Linq;
          
          static class WeakObservation
          {
              public static IObservable<T> ToWeakObservable<T>(this IObservable<T> observable)
              {
                  return Observable.Create<T>(observer =>
                      (IDisposable)new DisposableReference(new WeakObserver<T>(observable, observer), observer)
                      );
              }
          }
          
          class DisposableReference : IDisposable
          {
              public DisposableReference(IDisposable InnerDisposable, object Reference)
              {
                  this.InnerDisposable = InnerDisposable;
                  this.Reference = Reference;
              }
          
              private IDisposable InnerDisposable;
              private object Reference;
          
              public void Dispose()
              {
                  InnerDisposable.Dispose();
                  Reference = null;
              }
          }
          
          class WeakObserver<T> : IObserver<T>, IDisposable
          {
              private readonly WeakReference reference;
              private readonly IDisposable subscription;
              private bool disposed;
          
              public WeakObserver(IObservable<T> observable, IObserver<T> observer)
              {
                  this.reference = new WeakReference(observer);
                  this.subscription = observable.Subscribe(this);
              }
          
              public void OnCompleted()
              {
                  var observer = (IObserver<T>)this.reference.Target;
                  if (observer != null) observer.OnCompleted();
                  else this.Dispose();
              }
          
              public void OnError(Exception error)
              {
                  var observer = (IObserver<T>)this.reference.Target;
                  if (observer != null) observer.OnError(error);
                  else this.Dispose();
              }
          
              public void OnNext(T value)
              {
                  var observer = (IObserver<T>)this.reference.Target;
                  if (observer != null) observer.OnNext(value);
                  else this.Dispose();
              }
          
              public void Dispose()
              {
                  if (!this.disposed)
                  {
                      this.disposed = true;
                      this.subscription.Dispose();
                  }
              }
          }
          

          【讨论】:

          • 我刚刚测试了这段代码,它仍然泄漏。我强烈建议不要尝试这样做。整个思路都有问题。 1)静态重播主题 - 它永远不会释放它的缓存 2)如果你不为 Rx 实现 Dispose 模式,你还有什么不释放? - 事件处理程序、IO 连接? 3) 用户不能确定性地处置您的资源 4) 代码实际上不起作用 5) 少即是多。您有更多的代码会欺骗其他编码人员,让他们认为该代码有效,而实际上它不起作用,只会创建 100 行噪声代码。 请不要这样做
          • @Lee 回应 1、2、3:您可能会使用它的示例场景已被遗漏,您不会在需要确保源被处理的情况下使用它。这是为了确保听众得到处置。如果应用程序的整个生命周期都需要您的 Observable 是合适的,但是如果唯一引用它的东西就是那个 observable,则应该处置您的观察者。这几乎可以回答 4 - 即它确实可以确保观察者不被泄露
          • @Lee 回复 5:代码很多,而且它做了一些乍看之下似乎很简单,但结果却相对复杂的事情。编写这么一大段代码的理由是它具有高度的可重用性。只要您发现自己处于上述情况,这将起作用。它可以作为一个 API 公开,并带有一些适当的文档来说明弱可观察对象的作用,然后可以在无需了解其工作原理的情况下使用。
          • 我有兴趣看到这个工作的例子。我了解无法处置/收集来源。它是静态的,所以它永远存在。你如何建议订阅(即观察者)得到处置。我假设您实际上并没有创建和传递观察者,并且您正在使用订阅扩展方法来采取行动并为您隐式创建一个内部观察者,该观察者在完成/错误/取消订阅时被处置。
          • 我并不是想变得迟钝,只是当我运行你的代码并杀死窗口时,确保列表框停止填充但仍然引用观察者,订阅仍然运行并且它还是泄露了。但是,如果您只是意味着当您取消订阅/OnComplete/OnError 时将释放内部 observable,那么这已经发生了。当您使用订阅扩展方法时,这是 Rx 的默认行为。我希望我能接近理解。
          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2017-04-10
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2015-05-01
          相关资源
          最近更新 更多