【问题标题】:Using Reactive Extensions - is this async?使用响应式扩展 - 这是异步的吗?
【发布时间】:2013-09-24 01:57:09
【问题描述】:

尝试将现有数据访问代码转换为异步并遇到 Rx,因为您无法在方法主体中返回带有 yield returnTask<IEnumerable<T>>

我写了这个,但不确定它是异步的,所以感激地收到了指针

public class EmployeeRepository : IEmployeeRepository
{
    public IAsyncEnumerable<Employee> GetEmployees()
    {
        return Enumerable().ToAsyncEnumerable();
    }

    private IEnumerable<Employee> Enumerable()
    {
        using (var connection = new SqlConnection(ConfigurationManager.ConnectionStrings["DBConnString"].ConnectionString))
        {
            connection.Open();
            using (var command = new SqlCommand(@"SELECT * FROM EMPLOYEES", connection))
            {
                using (var reader = command.ExecuteReader())
                {
                    while (reader.Read())
                    {
                        yield return
                            new Employee()
                                {
                                    Id = ReadField<int>(reader, "Id"),
                                    Name = ReadField<string>(reader, "Name")
                                };
                    }
                }
            }
        }
    }

    private static T ReadField<T>(IDataRecord reader, string fieldName)
    {
        var value = reader[fieldName];
        return value == DBNull.Value ? default(T) : (T)value;
    }
}

【问题讨论】:

  • 什么意思,“你不能返回Task&lt;IEnumerable&lt;T&gt;&gt;?当然可以。你的意思是你当前的接口没有在它的方法签名中使用它吗?将代码转换为异步是通常是一项重大更改,而且非常重要;您将不得不更改许多接口和方法签名。您甚至应该更改名称,因为习惯上在异步方法的末尾包含后缀 Async
  • 反应式框架与这个问题有什么关系?
  • 你不能用yield return返回Task&lt;IEnumerable&lt;T&gt;&gt;
  • @Enigmativity 我正在尝试找到一种方法来实现收益回报,我认为 Rx 可以给我。
  • @Jon - yield return 仅适用于 IEnumerable&lt;&gt;,不适用于 IObservable&lt;&gt;。如果您使用yield return,它仍然是同步的。如果你走 Rx 路线,虽然你可以轻松让它异步,但不是 yield return 然后。

标签: c# .net asynchronous system.reactive


【解决方案1】:

这不是异步的。 ToAsyncEnumerable 创建一个简单的适配器,在每次调用 MoveNext 时都会阻塞。返回这样的异步适配器是不好的做法,与执行Task.Run(() =&gt; BlockingMethod()) 相同。它向用户隐藏了实现效率低下的问题,如果他们知道存在这种情况,他们可能能够以更好的方式解决此问题。

IAsyncEnumerable 没有语言集成的产量功能,但可以模拟。我have code to do it,但公平警告这会产生一些开销:

IAsyncEnumerable<Employee> async = AsyncEnumerableEx.Create<Employee>(
                                                  async (y, cancellationToken) =>
{
    using (var connection = new SqlConnection(ConfigurationManager
                            .ConnectionStrings["DBConnString"].ConnectionString))
    {
        await connection.OpenAsync(cancellationToken);
        using (var command = new SqlCommand(@"SELECT * FROM EMPLOYEES",
                                            connection))
        {
            using (var reader = await
                                   command.ExecuteReaderAsync(cancellationToken))
            {
                while (await reader.ReadAsync(cancellationToken))
                {
                    await y.YieldReturn(new Employee()
                    {
                        Id = ReadField<int>(reader, "Id"),
                        Name = ReadField<string>(reader, "Name")
                    });
                }
            }
        }
    }
});

如果你想使用实际的 Rx,它内置了一个几乎相同的 Observable.Create 实用程序。由于减少了一些等待开销,它会更有效率。

IObservable<Employee> async = Observable.Create<Employee>(
                                                    async (obs, cancellationToken) =>
{
    using (var connection = new SqlConnection(ConfigurationManager
                            .ConnectionStrings["DBConnString"].ConnectionString))
    {
        await connection.OpenAsync(cancellationToken);
        using (var command = new SqlCommand(@"SELECT * FROM EMPLOYEES",
                                            connection))
        {
            using (var reader = await
                                   command.ExecuteReaderAsync(cancellationToken))
            {
                while (await reader.ReadAsync(cancellationToken))
                {
                    obs.OnNext(new Employee()
                    {
                        Id = ReadField<int>(reader, "Id"),
                        Name = ReadField<string>(reader, "Name")
                    });
                }
            }
        }
    }
});

【讨论】:

  • 我有这个编译,但它只返回 1 名员工而不是 IEnumerable:
  • 你是如何使用它的?如果应该有更多记录,我认为它没有任何理由只返回一条记录。
  • 您必须在 Create 上提供一个泛型类型,那应该是什么?它必须是一个 Task 对吗? GetEmployees 应该返回一个 Task 是异步的
  • 这里有一个要点。数据库中有 2 个项目。它只返回 1。gist.github.com/jchannon/6673506
  • 很抱歉,没有测试该代码。泛型应该是Employee。 IObservable 已经是异步的,通过它传递任务几乎总是不正确的。帖子已编辑。
【解决方案2】:

如果您想使用 Rx,请尝试以下操作:

public IObservable<Employee> GetEmployees()
{
    return Observable.Create<Employee>(o =>
        Observable.Using(() => new SqlConnection(ConfigurationManager
            .ConnectionStrings["DBConnString"].ConnectionString),
            connection =>
                Observable.Using(() =>
                {
                    connection.Open();
                    return new SqlCommand(
                        @"SELECT * FROM EMPLOYEES", connection);
                },
                    command =>
                        Observable.Using(() => command.ExecuteReader(),
                            reader =>
                                Observable.Generate(
                                    0,
                                    x => reader.Read(),
                                    x => x,
                                    x => new Employee()
                                {
                                    Id = ReadField<int>(reader, "Id"),
                                    Name = ReadField<string>(reader, "Name")
                                }, Scheduler.Default)))).Subscribe(o));
}

【讨论】:

  • 这是异步的吗?我没有看到 async/await 和 Tas 返回
  • @Jon - 通过Scheduler.Default 异步,而不是通过新的async/await 关键字。 OP 的原始问题中没有async 的暗示。
  • 这是异步的,因为它发生在后台线程上,但不具备真正的异步操作所具有的任何可伸缩性。
  • @CoryNelson - 它正在从数据库中连续读取。您认为它可以获得多少可扩展性?
  • 相当多,但这不是我的意思。这是我提到的“伪异步”作为不好的做法。您可以相当轻松地修改此示例以使用真正的异步。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-02-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-07-20
相关资源
最近更新 更多