【问题标题】:Better C# design to traverse file system asynchronously to process large number of files on daily basis更好的 C# 设计以异步遍历文件系统以每天处理大量文件
【发布时间】:2019-07-28 21:22:21
【问题描述】:

我正在创建一个 c# 控制台应用程序,它将遍历给定文件夹(和子文件夹)以加密所有文件(二进制或文本)并更新 sqlserver 数据库中的IsEncrypted 标志。客户端上将有数百万个文件需要加密。我们计划在每天的非工作时间(例如从每晚 10 点开始运行 8 小时)按计划任务运行应用程序。

我有两个选择:

选项 1

使用Parallel.ForEach 处理文件。

public void Process(ProcessorOptions options, ProcessorParameter parameter)
{
    int counter = 0;
    CancellationTokenSource cts = new CancellationTokenSource();
    ParallelOptions parallelOptions = new ParallelOptions();
    parallelOptions.CancellationToken = cts.Token;

    try
    {
        parallelOptions.MaxDegreeOfParallelism = Environment.ProcessorCount;
        if (options.NumberOfThreads > 0)
        {
            parallelOptions.MaxDegreeOfParallelism = options.NumberOfThreads;
        }

        if (options.StopTime != 0)
        {
            Timer timer = new Timer(callback => { cts.Cancel(); }, null, options.StopTime * 60000, Timeout.Infinite);
        }

        List<string> storagePaths = parameter.StoragePaths;
        Log("Process Started...");

        foreach (var path in storagePaths)
        {
            Parallel.ForEach(TraverseDirectory(path, f => f.Extension != ".enc"), parallelOptions, file =>
            {
                if (file.Name.IndexOf("SRSCreate.dir") < 0)
                {
                    ProcessFile(parameter, file.FullName, file.Directory.Name, file.Name);
                    counter++;
                }
            });
        }
        Log(string.Format("Process Files Ended... Total File Count = {0}", counter));
    }
    catch (OperationCanceledException ex)
    {
        log.WriteWarningEntry(string.Format("Reached stop time = {0} min, explicit cancellation triggered. Total number of files processed = {1}", options.StopTime, counter.ToString()), ex);
    }
    catch (Exception ex)
    {                
        log.WriteErrorEntry(ex);
    }
    finally
    {
        cts.Dispose();
    }
}

我做了基准测试,发现处理 2000 个文件几乎需要 7-8 分钟。我可以做些什么来提高性能吗?此外,确定下一次运行(第二天)从哪里开始的最佳方法是什么?

选项 2

使用RabbitMQ 的现有设计来推送带有文件路径的消息,以处理文件以实现可扩展性和维护列表。

public void Process(ProcessorOptions options, ProcessorParameter parameter)
{
    try
    {
        using (IConnection connection = parameter.ConnectionFactory.CreateConnection())
        {
            using (IModel channel = connection.CreateModel())
            {
                var queueName = parameter.TopicSubscription.DeriveQueueName();
                var queueDeclareResponse = channel.QueueDeclare(queueName, true, false, false, null);
                EventingBasicConsumer consumer = new EventingBasicConsumer(channel);

                consumer.Received += (o, e) =>
                {
                    string messageContent = Encoding.UTF8.GetString(e.Body);
                    FileData message = JsonConvert.DeserializeObject(messageContent, typeof(FileData)) as FileData;
                    ProcessFile(parameter, message.EntityId, message.Attributes["Id"], message.Attributes["filename"]);
                };

                string consumerTag = channel.BasicConsume(queueName, true, consumer);
            }
        }
    }
    catch (Exception ex)
    {
        log.WriteErrorEntry(ex);
    }
    finally
    {
        Trace.Exit(method);
    }
}

在配置StopTime 之后,我仍然需要弄清楚如何停止阅读消息。性能不是很好,我看到处理 2000 个文件大约需要 25 - 30 分钟。我们认为我们可以在一台机器或多台机器上运行应用程序的多个副本来处理单个队列以进行扩展。您认为,我可以更改此代码以使其更优化吗?

最后一个问题:您认为是否还有其他选项比上述选项更高效和可扩展?

注意:

1) 方法ProcessFile 调用加密逻辑和更新数据库的逻辑。

2)我们遍历文件夹而不是从数据库开始,因为文件系统中可能存在数据库中尚不存在的文件。

【问题讨论】:

  • 你很可能会达到盒子的 io 限制。首先检查实际的限制组件。否则这很难回答。哦,一个问题:如果您的数据库不同步怎么办?这会是个大问题吗?
  • 这个问题有点问题,因为你在一个问题中问了很多事情......这一切都归结为瓶颈是什么,磁盘是什么(ssd?你在写回相同的驱动器?还是不同的驱动器?)。使用单线程读取器、并行加密器、单线程写入器可能会更好
  • @Stefan,这些文件存储在服务器上,由客户端应用程序查看。客户端应用程序依靠数据库根据 IsEncrypted 标志来识别文件是否需要解密才能查看。所以,数据库必须同步。
  • @KeithNicholas,很抱歉在一篇文章中问了这么多问题。我只想把所有东西都放在一个地方,我相信这个问题也会对其他人有所帮助。我同意有时你可以用代码做很多事情,然后你必须考虑硬件。这就是为什么我们考虑 RabbitMQ 的可扩展性的原因,如果客户端有硬件,他们可以投入更多的计算机、磁盘 (SSD) 等来加倍处理。单线程读写器不会阻塞系统?
  • 好吧,您需要最大化磁盘吞吐量,从多个线程访问磁盘可能会破坏磁盘的任何读/写缓存。我会做一堆实验来看看。我真的不明白你认为rabbitmq会为你做什么。但是您根本没有真正描述过部署架构。从您的问题看来,磁盘-> 内存-> 加密-> 磁盘....如果加密比磁盘内存更昂贵并再次返回磁盘,那么使用分布式系统进行加密可能会有所帮助。但我可能会为此使用 Akka 之类的东西。

标签: c# rabbitmq filesystems


【解决方案1】:

这属于性能问题的范畴,所以我将首先链接性能咆哮:https://ericlippert.com/2012/12/17/performance-rant/

这个操作本质上应该是 Diskbound,而不是 CPU bound。进程迭代文件的速度以及读取、加密和写入文件的速度有多快 - 都显然是磁盘绑定的。在磁盘上同时进行更多操作会使速度变慢,而不是变快。当然,除非你有一些极端的设置,比如 SSD 的 Raid 0。

如果有什么可以从多任务处理中受益,那应该是数据库访问。通常那些会通过网络堆栈,特别是如果数据库在另一台计算机上,它很有可能会比磁盘慢。同时,您不想通过查询向数据库发送垃圾邮件。所有查询都有开销,1 200 行查询比 200 1 行查询快。因此,以某种形式的枚举或流式方法获取数据库数据,然后遍历文件。但是哪一个真正最慢取决于每次运行时有多少新的/未加密的文件。

将整个东西移入数据库是可行的。有两种方法可以将 BLOBS 与 DB 一起存储,听起来您正在使用“存储在磁盘上,仅在 DB 中链接”。如果是这样,像 Filestream 这样的属性可能会对您有所帮助:https://www.red-gate.com/simple-talk/sql/learn-sql-server/an-introduction-to-sql-server-filestream/

有点偏离主题,但我的一个 Pet-Peeve 是异常处理,你的示例代码中有一个大罪:

catch (Exception ex)
{
    log.WriteErrorEntry(ex);
}

你捕捉到Exception 但不要让它继续,这意味着你在致命异常之后继续。那只会给你更多——更难理解的——后续例外。所以你永远不应该那样做。有两篇关于异常处理的文章,我确实链接了很多,我认为它们可能对您有所帮助:

【讨论】:

  • 很好的答案,但为什么你必须以the rant 开头?在这种情况下,性能显然是一个特性,而不是事后才想到的。
  • 我已经阅读了咆哮。如果您必须加密数百万个文件,并且在非工作时间安排日常任务,那么您就有一个性能问题需要解决(而且没什么好抱怨的)。
  • @TheodorZoulias 请重新阅读第 2、3 和 4 部分。如果您仍然需要帮助以了解它在此处的应用方式,我可以向您详细解释。
  • 我不会拒绝一个详细的解释,因为你愿意提供一个,关于这里如何适用@EricLippert 的rant 的以下几点:2)你真的需要回答这个问题? 3)这真的是瓶颈吗? 4) 差异是否相关?
  • @TheodorZoulias 2) 处理问题是否存在存在需要解决的问题。如果代码对用户来说已经足够快了,那么这项工作只是一个毫无意义的改变,有很多危险,没有任何收获。任何人都可以想到的唯一可能的改进是添加 Multtiasking。如果你做错了,你最终会得到更复杂、更容易出错、需要更多内存并且比你开始时慢的代码。 “除非发现性能问题,否则不要进行基于性能的更改” |操作没有提到实际问题。
【解决方案2】:

我不确定生产中会涉及多少物理驱动器。但如果需要,客户可以添加更多。未加密的文件被同一服务器上的加密文件替换,100%的文件需要加密,因为未加密的文件存在安全风险,并且每天,计数都会下降。是的,加密要求文件在内存中才能运行算法。文件的平均大小约为 3 mb。我知道文件大小没有限制,但通常我们会得到巨大的图像文件、word 和 excel doc,然后是一些小的文本文件。

我可以看到问题有很多未知数,这表明单个配置不会在所有情况下都解决问题。所以我的建议是使系统灵活。我将从为所涉及的物理驱动器创建配置开始。每个物理驱动器都应该有一个并发设置。 SSD 驱动器可以在 2-3 个线程同时读取或写入的情况下以最佳方式工作,而硬盘驱动器可能会遇到多个线程。下一个重要的设置是加密线程的数量。理想情况下,当正在运行的线程数等于机器的可用处理器/内核数时,系统应该工作得最好。执行流程应该是这样的:

  • IO 线程正在从/向其关联的物理驱动器读取或写入文件。
  • 当 IO 线程完成对未加密文件的读取后,会将其排入全局队列以进行处理。
  • 当 IO 线程完成写入加密文件时,它还会更新数据库。
  • 加密器线程不断汇集全局队列以供文件处理。
  • 当加密线程完成对文件的处理后,它会将其排入文件物理驱动器的专用队列中。
  • 当 IO 线程空闲时,如果有任何已处理的文件,它会查看其关联物理驱动器的专用队列。如果有,它会将其出列并将其写入磁盘。如果没有,它会继续从磁盘读取另一个文件。

所有这些都可以通过线程或任务以及BlockingCollection 类来实现。不需要Parallel.ForEach 或第三方库。

【讨论】:

  • 我可能误读了您关于物理驱动器数量的问题。您是指未加密文件所在的物理驱动器数量吗?如果是,则所有未加密的文件都存储在一个驱动器中的一个位置。有一个文件夹,它由以文档 ID 的 2 位前缀命名的子文件夹进一步组织。你说的全局队列,是指像rabbitmq这样的物理队列吗?
  • 如果所有文件都存储在一个驱动器中,那么事情会简单得多,但性能改进的机会也会大大减少。单个驱动器的 IO 速度几乎肯定会成为该过程的限制因素(瓶颈)。队列是指ConcurrentQueue,包裹在BlockingCollection 中。老实说,我对 RabbitMQ 一无所知。
  • 顺便说一句,如果未加密文件与加密文件存储在不同的物理驱动器中,那么立即有机会提升 x2 性能。
猜你喜欢
  • 2018-07-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多