【问题标题】:Enforce Observable Subscribers to only write to the stream one at a time强制 Observable 订阅者一次只写入一个流
【发布时间】:2015-09-12 21:13:16
【问题描述】:

我目前正在使用 observables 来管理总线上生成的消息,这些消息被推送到各种流中。

一切正常,但由于消息可以进入,系统可能会尝试一次将多条消息写入流(即来自多个线程的消息),或者消息的发布速度比写入速度快流...如您所见,这会在写入时引起问题。

因此,我试图弄清楚如何组织事物,以便在收到消息时一次只处理一个。有什么想法吗?

public class MessageStreamResource : IResourceStartup
{
    private readonly IBus _bus;
    private readonly ISubject<string> _sender;

    public MessageStreamResource(IBus bus)
    {
        _bus = bus;

        _senderSubject = new Subject<string>();

        //`All` can publish messages at the same time as it's
        //collecting data being generated from different threads
        _bus.All.Subscribe(message => Observable.Start(() => ProcessMessage(message), TaskPoolScheduler.Default));

        //Note the above hops off the calls context so that the 
        //writing to the stream wont slow down the caller.
    }

    public void Configure(IAppBuilder app)
    {
        app.Map("/stream", async context =>
        {
            ...

            await context.Response.WriteAsync("Lets party!\n");
            await context.Response.Body.FlushAsync();

            var unSubscribe = _sender.Subscribe(async t =>
            {
                //PROBLEM HERE
                //I only want this callback to be executed 
                //one at a time...

                await context.Response.WriteAsync($"{t}\n");
                await context.Response.Body.FlushAsync();
            });

            ...

            await HoldOpenTask;
        });
    }

    private void ProcessMessage(IMessage message)
    {
        _sender.OnNext(message.Payload);
    }
}

【问题讨论】:

  • 如果您使用 Rx,那么它会自动确保一次处理一条消息。这就是 Rx 为您所做的。
  • 使用_sender.Synchronize().Subscribe(...)有问题吗?

标签: c# .net task-parallel-library system.reactive observable


【解决方案1】:

如果我正确理解了这个问题,这可能可以通过SemaphoreSlim 来完成:

// ...
var semaphore = new SemaphoreSlim(initialCount: 1);

var unSubscribe = _sender.Subscribe(async t =>
{
    //PROBLEM HERE
    //I only want this callback to be executed 
    //one at a time...

    await semaphore.WaitAsync();
    try
    {
        await context.Response.WriteAsync($"{t}\n");
        await context.Response.Body.FlushAsync();
    }
    finally
    {
        semaphore.Release();
    }
});

SemaphoreSlimIDisposable,请务必在适当的时候将其丢弃。

更新,第二次看,MapExtensions.Map 接受Action&lt;IAppBuilder&gt;,所以你传递了一个async void lambda,本质上创建了一堆即发即弃的异步操作。 Map 呼叫将返回给呼叫者,而他们可能仍在徘徊。这很可能不是您想要的,是吗?

【讨论】:

  • 所以在我发布问题后,我添加了一个SemaphoreSlim,但希望有一种“原生”的方式来处理 observables。
  • 另外,MapExtensions.Map 我的代码比列出的要复杂一些,但我试图让它更简单一些。
  • @anthonyv,如果你不喜欢SemaphoreSlim,你可以创建一个自定义任务调度程序或使用EventLoopScheduler(如果你没有任何阻塞代码应该没问题)。我对async void lambdas 的担忧仍然有效。如果您发布更多代码可能会有所帮助,至少可以显示您如何调用app.Run
  • 感谢您的帮助!
猜你喜欢
  • 2023-03-15
  • 1970-01-01
  • 1970-01-01
  • 2019-12-03
  • 2020-08-16
  • 2019-05-10
  • 2018-01-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多