【问题标题】:Async Producer/Consumer异步生产者/消费者
【发布时间】:2013-09-11 18:59:43
【问题描述】:

我有一个从多个线程访问的类的实例。此类接受此调用并将元组添加到数据库中。我需要以串行方式完成此操作,因为由于某些 db 限制,并行线程可能会导致数据库不一致。

由于我不熟悉 C# 中的并行性和并发性,我这样做了:

private BlockingCollection<Task> _tasks = new BlockingCollection<Task>();

public void AddDData(string info)
{
    Task t = new Task(() => { InsertDataIntoBase(info); });
    _tasks.Add(t);
}

private void InsertWorker()
{
    Task.Factory.StartNew(() =>
    {
        while (!_tasks.IsCompleted)
        {
            Task t;
            if (_tasks.TryTake(out t))
            {
                t.Start();
                t.Wait();
            }
        }
    });
}

AddDData 是由多个线程调用的,InsertDataIntoBase 是一个非常简单的插入,应该需要几毫秒。

问题是,由于某种原因,我缺乏知识并无法弄清楚,有时一个任务被调用了两次!它总是这样:

T1 T2 T3 T1

我是否理解 .Take() 完全错误,是我遗漏了什么还是我的生产者/消费者实现真的很糟糕?

最好的问候, 拉斐尔

更新:

按照建议,我用这种架构做了一个快速的沙盒测试实现,正如我所怀疑的,它不能保证在前一个任务完成之前不会触发任务。

所以问题仍然存在:如何正确地对任务进行排队并按顺序触发它们?

更新 2:

我简化了代码:

private BlockingCollection<Data> _tasks = new BlockingCollection<Data>();

public void AddDData(Data info)
{
    _tasks.Add(info);
}

private void InsertWorker()
{
    Task.Factory.StartNew(() =>
    {
        while (!_tasks.IsCompleted)
        {
            Data info;
            if (_tasks.TryTake(out info))
            {
                InsertIntoDB(info);
            }
        }
    });
}

请注意,我摆脱了 Tasks,因为我依赖于同步的 InsertIntoDB 调用(因为它在循环内),但仍然没有运气......一代很好,我绝对确定只有唯一的实例是去排队。但不管我怎么尝试,有时同一个对象会被使用两次。

【问题讨论】:

  • 你是如何生成主键的?
  • 实际上我简化了这里显示的代码,因为数据不是字符串,而是一个非常复杂的对象。 PK 实际上是 2 个对象字段(名称字符串和日期时间值)。我无法控制数据库。
  • 我认为一个简单的lock 足以序列化调用。
  • 要确认您确实在执行两次相同的任务,请将 PK 和 TaskID (Task.CurrentId) 写入命令行并查看输出。启动 1M+ 任务时我无法重现此问题...
  • 来自msdn:删除项目的顺序取决于用于创建 BlockingCollection 实例的集合类型。创建 BlockingCollection 对象时,可以指定要使用的集合类型。例如,您可以为先进先出 (FIFO) 行为指定 ConcurrentQueue 对象。 BlockingCollection 的默认集合类型是 ConcurrentQueue。

标签: c# multithreading concurrency parallel-processing producer-consumer


【解决方案1】:

我认为这应该可行:

    private static BlockingCollection<string> _itemsToProcess = new BlockingCollection<string>();

    static void Main(string[] args)
    {
        InsertWorker();
        GenerateItems(10, 1000);
        _itemsToProcess.CompleteAdding();
    }

    private static void InsertWorker()
    {
        Task.Factory.StartNew(() =>
        {
            while (!_itemsToProcess.IsCompleted)
            {
                string t;
                if (_itemsToProcess.TryTake(out t))
                {
                    // Do whatever needs doing here
                    // Order should be guaranteed since BlockingCollection 
                    // uses a ConcurrentQueue as a backing store by default.
                    // http://msdn.microsoft.com/en-us/library/dd287184.aspx#remarksToggle
                    Console.WriteLine(t);
                }
            }
        });
    }

    private static void GenerateItems(int count, int maxDelayInMs)
    {
        Random r = new Random();
        string[] items = new string[count];

        for (int i = 0; i < count; i++)
        {
            items[i] = i.ToString();
        }

        // Simulate many threads adding items to the collection
        items
            .AsParallel()
            .WithDegreeOfParallelism(4)
            .WithExecutionMode(ParallelExecutionMode.ForceParallelism)
            .Select((x) =>
            {
                Thread.Sleep(r.Next(maxDelayInMs));
                _itemsToProcess.Add(x);
                return x;
            }).ToList();
    }

这确实意味着消费者是单线程的,但允许多个生产者线程。

【讨论】:

  • 我尝试了超过 10000000 次迭代而没有重复。会再考虑一下。
  • +1。请注意,您可以通过将 while (!IsCompleted) 和 TryTake 替换为单个 foreach (string t in _itemsToProcess.GetConsumingEnumerable()) 来简化您的工作人员。见GetConsumingEnumerable。
  • @JimMischel 我考虑过,但是这里的排序很重要,this msdn document 说使用这种方法时不能保证排序。
  • @sga101 经过几个小时的调试和记录,我发现错误不在这部分代码中,而是在数据库插入中。似乎有一个令人讨厌的错误是由我用这种代码气味创建的竞争条件引起的。当你的代码让我发现这一点时,为你竖起大拇指,万分感谢!
  • @sga101:我假设你说的是这样的行:“不能保证项目的枚举顺序与生产者线程添加它们的顺序相同。”老实说,我不知道那条线应该是什么意思。我向你保证GetConsumingEnumerable 确实 以先进先出的顺序移除东西。查看 .NET Framework 源代码可以确认:GetConsumingEnumerable 只是循环中的一堆 TryTake 调用。项目的删除顺序与插入的顺序相同。
【解决方案2】:

来自您的评论

“我简化了这里显示的代码,因为数据不是字符串”

我假设传递给 AddDData 的 info 参数是可变引用类型。确保调用者没有使用相同的 info 实例进行多次调用,因为该引用是在 Task lambda 中捕获的。

【讨论】:

  • 是的,它是一个可变引用类型,我会仔细检查它(尽管我认为这不是问题)。竞争条件可能会导致这种行为,但使用 BlockingCollection 作为缓冲区不会降低这种风险?
  • Rafa Borges:我误读了代码 - 没有竞争条件,因为它确保项目按顺序执行。不排除在不同任务中捕获相同物品的可能性。
【解决方案3】:

根据您提供的跟踪信息,唯一合乎逻辑的可能性是您调用了两次(或更多次)InsertWorker。因此,有两个后台线程等待项目出现在集合中,有时它们都设法抓取一个项目并开始执行它。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-04-12
    • 1970-01-01
    • 2021-04-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多