【问题标题】:Close unmanaged resources when Subscription end in Reactive Extensions订阅以反应式扩展结束时关闭非托管资源
【发布时间】:2015-04-27 07:50:33
【问题描述】:

我正在从 Rx 向网络写入数据。当然,我会在订阅结束时使用Finally 关闭我的流。这在OnError()OnComplete() 上都可以正常工作。 Rx 将依次运行OnNext() ... OnNext()OnComplete()Finally()

但是,有时我想提前终止序列,为此我使用Dispose()。不知何故,Finally() 现在与最后一个OnNext() 调用并行运行,导致在OnNext() 中仍然写入流时出现异常,以及写入不完整。

我的订阅大致如下:

NetworkStream stm = client.GetStream();
IDisposable disp = obs
    .Finally(() => {
        client.Close();
    })
    .Subscribe(d => {
        client.GetStream().Write(d.a, 0, d.a.Lenght);
        client.GetStream().Write(d.b, 0, d.b.Lenght);
    } () => {
        client.GetStream().Write(something(), 0, 1);
    });
Thread.sleep(1000);
disp.Dispose();

我也尝试了替代方法,CancellationToken

如何正确取消订阅?我不介意它是否跳过OnComplete(),只要Finally() 仍在运行。但是,并行运行Finally() 是有问题的。

我也觉得应该有一种更好的方法来管理资源,将声明移到序列中,这将是一个更好的解决方案。

编辑:以下代码更清楚地显示了问题。我希望它总是打印 true,相反,它经常给出 false,表明 Dispose 在最后一个 OnNext 之前结束。

using System;
using System.Collections.Generic;
using System.Linq;
using System.Net.Sockets;
using System.Reactive;
using System.Reactive.Disposables;
using System.Reactive.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;

namespace ConsoleApplication1
{
    class Program
    {
        static void Main(string[] args)
        {
            Console.WriteLine("Try finally");
            for (int i = 0; i < 10; i++)
            {
                Finally();
            }
            Console.WriteLine("Try using");
            for (int i = 0; i < 10; i++)
            {
                Using();
            }
            Console.WriteLine("Try using2");
            for (int i = 0; i < 10; i++)
            {
                Using2();
            }
            Console.ReadKey();
        }

        private static void Using2()
        {
            bool b = true, c = true, d;
            var dis = Disposable.Create(() => c = b);
            IDisposable obDis = Observable.Using(
                () => dis,
                _ => Observable.Create<Unit>(obs=>
                    Observable.Generate(0,
                    i => i < 1000,
                    i => i + 1,
                    i => i,
                    i => TimeSpan.FromMilliseconds(1)
                ).Subscribe(__ => { b = false; Thread.Sleep(100); b = true; })))
                .Subscribe();
            Thread.Sleep(15);
            obDis.Dispose();
            d = b;
            Thread.Sleep(101);
            Console.WriteLine("OnDispose: {1,5} After: {2,5} Sleep: {0,5}", b, c, d);
        }

        private static void Using()
        {
            bool b = true, c = true, d;
            var dis = Disposable.Create(() => c = b);
            IDisposable obDis = Observable.Using(
                () => dis,
                _ => Observable.Generate(0,
                    i => i < 1000,
                    i => i + 1,
                    i => i,
                    i => TimeSpan.FromMilliseconds(1)
                )).Subscribe(_ => { b = false; Thread.Sleep(100); b = true; });
            Thread.Sleep(15);
            obDis.Dispose();
            d = b;
            Thread.Sleep(101);
            Console.WriteLine("OnDispose: {1,5} After: {2,5} Sleep: {0,5}", b, c, d);
        }

        private static void Finally()
        {
            bool b = true, c = true, d;
            IDisposable obDis = Observable.Generate(0,
                i => i < 1000,
                i => i + 1,
                i => i,
                _ => DateTime.Now.AddMilliseconds(1)
                )
                .Finally(() => c = b)
                .Subscribe(_ => { b = false; Thread.Sleep(100); b = true; });
            Thread.Sleep(15);
            obDis.Dispose();
            d = b;
            Thread.Sleep(101);
            Console.WriteLine("OnDispose: {1,5} After: {2,5} Sleep: {0,5}", b, c, d);
        }
    }
}

【问题讨论】:

    标签: c# system.reactive


    【解决方案1】:

    Finally 很可能不是您想要的。当您取消订阅时,它不会释放您的资源。相反,它的行为类似于 C# 中的普通 finally 块,也就是说,无论其对应的 try-块中的代码是否抛出异常,它都会保证执行某些代码。此外,鉴于this question on MSDN,您在Finally 中的代码甚至可能不会在任何情况下都执行,因为您的订阅未指定错误处理程序。

    你可能想要的是Using:

    IDisposable disp = Observable
        .Using(
            () => Disposable.Create(() => client.Close),
            _ => obs)
        .Subscribe(....);
    

    Using 负责在可观察对象终止或订阅被取消时正确处置资源。

    假设client是一个TcpClient,那就更简单了:

    IDisposable disp = Observable
        .Using(
            () => client),
            _ => obs)
        .Subscribe(....);
    

    我希望对OnNext 的调用不会与关闭客户端重叠,即使提早取消订阅也是如此,但我尚未对此进行测试。

    最后一件事:注意在您的示例中关闭外部变量,例如 stm。始终与当地人合作更安全。我会尝试的完整重写是这样的:

    IDisposable disp = Observable.Using(
        () => client,
        _ => Observable.Using(
             () => client.GetStream(),
             stream => Observable.Create<Unit>(observer => obs
                 .Subscribe(
                     d => {
                         stream.Write(d.a, 0, d.a.Lenght);
                         stream.Write(d.b, 0, d.b.Lenght);
                     },
                     () => {
                         stream.Write(something(), 0, 1);
                     }))))
        .Subscribe();
    

    【讨论】:

    • 实际上 Finally 确实会出现错误,我认为该错误自 2012 年以来已修复。无论如何,您的答案的重置看起来很有希望,我觉得 Observable.Using 应该适合,但我不能还不知道怎么做。 client 确实是 TcpClient,所以我要试一试。
    • 顺便说一句,您的最后一个示例无法编译。第二行有括号错误,Observable.Create 有问题我无法修复。
    • 应该修复。再试一次。
    • OnNext 调用重叠。这是 Rx 行为契约的一部分。
    • @Enigmativity 是的,OnNextOnComplete 不重叠。 Finally 然而,完全出乎意料,确实如此。 (但这并不违反任何 Rx 合同。)
    【解决方案2】:

    我认为您只是对NetworkStream 的运作方式做出了错误的假设。

    NetworkStream.WriteTcpClient.Close 不一定要等待客户端读取数据。 (此外,NetworkStream.Flush 什么也不做)。

    当您调用Close 时,您可能在客户端读取所有内容之前关闭了套接字。

    看看这个相关的问题:NetworkStream doesn't always send data unless I Thread.Sleep() before closing the NetworkStream. What am I doing wrong?

    该页面不同地提到了诸如使用接受超时的Close 的重载或指定LingerOption - 但最好是发送Shutdown 或具有更高级别的消息传递抽象,其中客户端确认您的带有自己回复的消息。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-12-27
      • 2023-03-04
      • 1970-01-01
      • 2013-05-15
      • 1970-01-01
      • 2013-02-02
      • 1970-01-01
      相关资源
      最近更新 更多