【发布时间】:2009-11-20 07:57:43
【问题描述】:
响应式扩展提供了许多帮助方法,用于将现有事件和异步操作转换为可观察对象,但是您将如何从头开始实现 IObservable
IEnumerable 具有可爱的 yield 关键字,使其实现起来非常简单。
实现 IObservable
我需要担心线程安全吗?
我知道支持在特定的同步上下文中回调,但这是我作为 IObservable
更新:
这是我的 C# 版本的 Brian 的 F# 解决方案
using System;
using System.Linq;
using Microsoft.FSharp.Collections;
namespace Jesperll
{
class Observable<T> : IObservable<T>, IDisposable where T : EventArgs
{
private FSharpMap<int, IObserver<T>> subscribers =
FSharpMap<int, IObserver<T>>.Empty;
private readonly object thisLock = new object();
private int key;
private bool isDisposed;
public void Dispose()
{
Dispose(true);
}
protected virtual void Dispose(bool disposing)
{
if (disposing && !isDisposed)
{
OnCompleted();
isDisposed = true;
}
}
protected void OnNext(T value)
{
if (isDisposed)
{
throw new ObjectDisposedException("Observable<T>");
}
foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
{
observer.OnNext(value);
}
}
protected void OnError(Exception exception)
{
if (isDisposed)
{
throw new ObjectDisposedException("Observable<T>");
}
if (exception == null)
{
throw new ArgumentNullException("exception");
}
foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
{
observer.OnError(exception);
}
}
protected void OnCompleted()
{
if (isDisposed)
{
throw new ObjectDisposedException("Observable<T>");
}
foreach (IObserver<T> observer in subscribers.Select(kv => kv.Value))
{
observer.OnCompleted();
}
}
public IDisposable Subscribe(IObserver<T> observer)
{
if (observer == null)
{
throw new ArgumentNullException("observer");
}
lock (thisLock)
{
int k = key++;
subscribers = subscribers.Add(k, observer);
return new AnonymousDisposable(() =>
{
lock (thisLock)
{
subscribers = subscribers.Remove(k);
}
});
}
}
}
class AnonymousDisposable : IDisposable
{
Action dispose;
public AnonymousDisposable(Action dispose)
{
this.dispose = dispose;
}
public void Dispose()
{
dispose();
}
}
}
编辑:如果 Dispose 被调用两次,不要抛出 ObjectDisposedException
【问题讨论】:
-
Wes Dyer 现在在 Channel9 上有一个视频,讨论这些接口的合同。
-
(30 秒后...channel9.msdn.com/posts/J.Van.Gogh/…)
-
酷 - 一定会看的:)
-
当 Dispose() 被调用两次时,你不应该抛出 ObjectDisposed。在调用 dispose 后调用其他方法时,您应该抛出 ObjectDisposed。
-
@JohnGietzen 你是对的。我编辑了代码以反映这一点。
标签: c# system.reactive