You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Rx.NET中如何正确处置Observable.Create创建的可释放资源?

Rx.NET 中正确处置Observable.Create发射的IDisposable资源

我在项目中使用Rx(.NET)时,遇到了一个关键问题:如何正确处置在Observable.Create()期间创建、并通过OnNext()发射的IDisposable资源。这里的RunData包含数据库事务和Autofac LifetimeScope,必须在订阅者完成数据处理后释放。

现有实现代码

无限序列生成逻辑

var obs = Observable.Create<RunData>(async (o) =>
{
    if (someCondition)
    {
        RunData runData = await CreateRunData(); // RunData是IDisposable,需要释放
        o.OnNext(runData);
    }

    o.OnCompleted();

    return Disposable.Empty;
})
.Concat(Observable.Empty<RunData>().Delay(TimeSpan.FromSeconds(2)))
.Repeat() // 序列完成后无限重订阅
.Publish().RefCount(); // 共享订阅

序列转换逻辑

var final = obs
    .Select(runData => /* 业务转换逻辑 */)
    // 其他操作符处理
    .Select(tuple => (tuple.runData, tuple.result));

订阅处置逻辑

final.Subscribe(
    async (tuple) =>
    {
        var (runData, result) = tuple;
        try
        {
            // 使用runData和result执行异步业务逻辑
        } catch (Exception e)
        {
            // 错误处理
        } finally 
        {
            // 手动释放runData
            runData.Dispose();
        }
    },
    (Exception e) =>
    {
        // 错误处理
    });

现有实现的问题

  • 泄漏风险:若序列中途抛出异常、或多订阅者场景下,RunData可能无法被正确处置
  • 职责错位:让订阅者负责资源处置不符合单一职责原则,增加维护成本
  • 尝试过的方案局限性:
    • Observable.Using():仅在序列结束时释放资源,但我的序列是无限的(需要用Scan()等操作构建中间状态),无法适用
    • Observable.Create()返回的Disposable回调:会在Create逻辑完成后立即触发,而非订阅者处理完数据后,引发竞态条件

正确实现方案

核心思路

让资源生命周期与单个数据项的处理周期绑定,将资源释放逻辑封装在序列内部,订阅者仅负责业务处理,无需关心资源管理。

方案1:用信号跟踪处理完成状态(适合异步处理场景)

修改序列生成逻辑,为每个RunData添加处理完成信号,确保资源在所有订阅者处理完成后释放:

var obs = Observable.Create<(RunData, IObserver<Unit>)>(async (observer, cancellationToken) =>
{
    var disposables = new CompositeDisposable();

    // 构建2秒周期的任务流
    var periodicStream = Observable.Interval(TimeSpan.FromSeconds(2))
        .SelectMany(async _ =>
        {
            if (someCondition)
            {
                var runData = await CreateRunData();
                // 创建异步信号,跟踪处理完成状态
                var processedSignal = new AsyncSubject<Unit>();
                
                // 发射资源+处理信号
                observer.OnNext((runData, processedSignal.AsObserver()));

                // 等待处理完成后释放资源
                await processedSignal.FirstAsync(cancellationToken);
                runData.Dispose();
            }
            return Unit.Default;
        })
        .Subscribe(_ => { }, observer.OnError, observer.OnCompleted);

    disposables.Add(periodicStream);
    return disposables;
})
// 保留信号的同时执行业务转换
.Select(item =>
{
    var (runData, signal) = item;
    // 你的业务转换逻辑,生成result
    var result = /* ... */;
    return (runData, result, signal);
})
.Publish().RefCount();

订阅时,处理完成后触发信号:

final.Subscribe(
    async tuple =>
    {
        var (runData, result, processedSignal) = tuple;
        try
        {
            // 执行异步业务逻辑
            await ProcessResult(runData, result);
        }
        catch (Exception e)
        {
            // 错误处理
        }
        finally
        {
            // 触发处理完成信号,通知序列释放资源
            processedSignal.OnNext(Unit.Default);
            processedSignal.OnCompleted();
        }
    },
    ex => /* 错误处理 */
);

方案2:Using+Defer+Repeat(适合同步/短周期处理场景)

如果每个RunData仅需单次处理,可将Observable.Using与Defer、Repeat结合,让每个周期创建独立的资源上下文,确保资源在单次数据发射后(处理完成前)不被释放:

var obs = Observable.Defer(() =>
    Observable.Using(
        // 每个周期创建新资源
        async () => await CreateRunData(),
        runData =>
        {
            if (someCondition)
            {
                // 发射资源后立即完成当前子序列,触发Using释放资源
                return Observable.Return(runData);
            }
            return Observable.Empty<RunData>();
        }
    )
)
// 等待2秒后重复
.Concat(Observable.Empty<RunData>().Delay(TimeSpan.FromSeconds(2)))
.Repeat()
.Publish().RefCount();

注意:该方案仅适用于订阅者处理逻辑为同步、或Rx操作符能保证处理完成后才触发序列完成的场景。若存在异步处理,仍需方案1的信号跟踪机制。

多订阅者场景补充

使用Publish().RefCount()时,多个订阅者会共享同一个RunData,此时需要跟踪所有订阅者的处理完成状态,可将AsyncSubject替换为CountDownEvent或使用Observable.WhenAll合并所有订阅者的处理信号,确保所有订阅者处理完成后再释放资源。


内容的提问来源于stack exchange,提问作者Maximilian Csuk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 02:10:49