【问题标题】:Interleaving processes on two threads在两个线程上交错进程
【发布时间】:2015-06-16 19:27:32
【问题描述】:

我有一个可以用来访问一些表格数据的库。这个库是我访问数据的唯一方法。我使用的方法需要一个查询字符串和一个为每个结果行调用的回调。

目前,回调将每一行加载到一个列表中,然后返回该列表。我想使用迭代器模式,但我对数据的唯一访问是通过这个回调方法。

有没有办法可以在第二个线程上运行查询/回调并将该代码与迭代器代码交错? 伪代码:

IEnumerable<Row> QueryData(string queryString)
{
    var callerLock = create new sync lock;
    var callbackLock = create new sync lock;
    var rows = create new stack of rows with capacity 1;
    var qthread = create new thread with QueryCallback(queryString, callerLock, callbackLock, rows);

    start qthread;
    while (qthread is running)
    {
        signal callbackLock;
        wait for callerLock;

        if stack is empty
            break;

        var row = pop from rows;
        yield return row;
    }
}

void QueryCallback(string queryString, lock callerLock, lock callbackLock, Stack<Row> rows)
{
    DoQueryWithCallback(queryString, row =>
    {
        wait for callbackLock;
        push row to rows;
        signal callerLock;
    });

    signal callerLock;
}

我尝试使用 .NET Framework 中提供的大多数锁来实现这一点,但它们都不起作用。我记得尝试过 Semaphore、SemaphoreSlim、AutoResetEvent、ManualResetEvent 和 Mutex。

P.S.:DoQueryWithCallback 来自图书馆。它是一个原生库(ILSpy/Reflector/etc 无法反编译它)。我想这个函数看起来像这样:

long DoQueryWithCallback(string queryString, Callback rowCallback)
{
    do some setup;

    Row row;
    while (next(out row))
            rowCallback(row);

    do some teardown;
}

【问题讨论】:

  • 你想达到什么目的?
  • 我想模拟一个惰性求值的迭代器。
  • 您希望通过这样做获得什么?如果底层数据源不支持它,那么除了多线程和锁定的大量开销之外,您还能得到什么?
  • @xxbbcc:也许在循环中处理每个项目而不是通过回调的可读性,但不需要一次加载整个集合?并不是说我不同意你的开销......
  • @DarkFalcon 您可能是对的,但归根结底,如果这是获取数据的唯一方法,那么回调代码仍然会出现在某个地方。除此之外的任何其他东西都是某种开销。我不确定我是否会打扰 - 对于相当程度的复杂性,收益是非常值得怀疑的。最后,所有的等待都会大大减慢进程。

标签: c# multithreading synchronization


【解决方案1】:

如果我正确理解了伪代码,您希望在后台线程上触发 fetch 操作并在结果进入时使用迭代器产生结果,而不是等待整个 fetch 操作完成后再返回。我会改变的几件事:

  • 如果要保持获取行的顺序,请使用队列而不是堆栈
  • 信号/阻塞只需要采用一种方式 - 线程 yield 返回行需要等待获取线程将项目添加到队列中。无需阻塞获取线程

下面是一个使用TaskConcurrentQueueAutoResetEvent 的简单示例:

public IEnumerable<Row> GetRows(string query)
{
    using (var resetEvent = new AutoResetEvent(false))
    {
        var rows = new ConcurrentQueue<Row>();

        var queryTask = Task.Run(() => DoQueryWithCallback(query, r =>
        {
            rows.Enqueue(r);
            resetEvent.Set();
        }));
        queryTask.ContinueWith(t => resetEvent.Set()); // This ensures that queryTask.IsCompleted will be true in the while loop below

        while (resetEvent.WaitOne() && !queryTask.IsCompleted)
        {
            Row row;
            while (rows.TryDequeue(out row))
                yield return row;
        }
    }
}

编辑

其实有更好的方法使用BlockingCollection

public IEnumerable<Row> GetRows(string query)
{
    using (var rows = new BlockingCollection<Row>())
    {
        Task.Run(() =>
        {
            DoQueryWithCallback(query, r => rows.Add(r));
            rows.CompleteAdding();
        });

        while (!rows.IsCompleted)
            yield return rows.Take();
    }
}

【讨论】:

  • 我将此标记为答案,因为这可能是无需等待整个提取完成即可立即获得结果的最佳方式。但我仍然想知道如何在两个不同的线程上交错代码。我尝试使用 AutoResetEvents 来实现这一点,但回调似乎多次进入锁定,只有一次调用 Set。
猜你喜欢
  • 1970-01-01
  • 2010-09-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-08-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多