【发布时间】:2014-05-30 00:13:18
【问题描述】:
我正在尝试编写一个可观察的辅助函数,将嵌套序列合并为一个序列。换句话说,签名看起来像这样:
public IObservable<string> CreateNested(
Func<IObservable<string>> createOuter,
Func<string, IObservable<string>> createInner);
一个细节是这些序列包装了服务调用,因此每个序列最多有一个项目。
所以我的第一次尝试成功了。但它在我看来是不必要的冗长,而且,Wait 的使用打破了可观察的模式,因为来自任一序列的错误都会引发异常,而不是传播回返回的序列:
public IObservable<string> CreateNested(Func<IObservable<string>> createOuter, Func<string, IObservable<string>> createInner)
{
return Observable.StartAsync(() =>
{
return Task.Factory.StartNew(() =>
{
string outerResult = createOuter().Wait();
var inner = createInner(outerResult);
return inner.Wait();
});
});
}
我的第二次尝试更好一些,但仍然使用Wait。
public IObservable<string> CreateNested(Func<IObservable<string>> createOuter, Func<string, IObservable<string>> createInner)
return createOuter().FirstOrDefaultAsync()
.Select(result => createInner(result).Wait());
}
如果我用另一个“FirstOrDefaultAsync()”替换上面的“Wait”,那么我得到IObservable<IObservable<string>>。有没有正确的方法来“合并”这两个序列?
编辑为了完整起见,我的测试如下(预期输出为“hello world”)。
public class Tester
{
public void Test()
{
CreateNested(CreateOuter, CreateInner).Subscribe(Console.WriteLine);
}
private IObservable<string> CreateOuter()
{
return Observable.Create<string>(observer =>
{
Task.Factory.StartNew(() =>
{
Thread.Sleep(1000);
observer.OnNext("hello");
observer.OnCompleted();
});
return new Action(() => { Console.WriteLine("Outer subscriber released"); });
});
}
private IObservable<string> CreateInner(string key)
{
return Observable.Create<string>(observer =>
{
Task.Factory.StartNew(() =>
{
Thread.Sleep(1000);
observer.OnNext(key + " world");
observer.OnCompleted();
});
return new Action(() => { Console.WriteLine("Inner subscriber released"); });
});
}
private IObservable<string> CreateNested(Func<IObservable<string>> createOuter, Func<string, IObservable<string>> createInner)
{
// TODO
}
}
【问题讨论】:
-
您尝试过 SelectMany 吗?
-
@phoog ha ...你摇滚,先生。如果您愿意,请继续发布作为答案。
标签: c# system.reactive