【问题标题】:True concurrency on a collection in Windows WF 4.5Windows WF 4.5 中集合的真正并发
【发布时间】:2015-08-13 14:12:52
【问题描述】:

我正在修改以前编码为同步运行所有内容的现有 Windows Workflow Foundation 项目。但是,随着数据集的增长,这需要进行更改以满足性能要求。

我有什么:

在工作流程中,我有一个父序列工作流程,其中包含一些基本工作流程,这些工作流程基本上设置了一些服务并准备它们运行。 然后我拥有工作流的大部分工作,它由一个 ForEach 工作流组成,该工作流对大约 15000 个项目的集合进行操作,每个项目大约需要 1-3 秒来处理(时间大约是 70% 的 CPU , 10% 的网络延迟,20% 的数据库查询/访问)。显然这需要 WAYYYY 太长时间。我需要将这个时间提高大约 5 倍(大约需要 5-6 小时,需要达到大约 1 小时)

德利马:

在这个项目之前我从未使用过 Windows 工作流,所以我非常不熟悉如何在集合上实现并发执行的简单实现。

想法:

我阅读了不同的工作流活动,并决定ParallelForEach 工作流活动可能是要走的路。我的想法是,我只需用 ParallelForEach 工作流活动切换我的 ForEach 工作流活动,并以 Parallel.Foreach() 在任务并行库中的工作方式实现并发。不幸的是,这似乎不是 ParallelForEach 工作流活动的实现方式。 ParallelForEach 工作流活动似乎只是将每个迭代放在一个堆栈中并几乎同步地对它们进行操作,除非工作流的主体是“空闲”(我认为这与 I/O 上的“等待”不同。这似乎是需要在工作流活动上设置的显式状态 - 每个 MSDN:

ParallelForEach 枚举其值并为 Body 调度 它枚举的每个值。它只安排身体。身体怎么样 执行取决于 Body 是否空闲。 如果身体不 闲置,它以相反的顺序执行,因为预定的 活动作为堆栈处理,最后一个计划的活动 首先执行。例如,如果您有一个集合 {1,2,3,4} 在 ParallelForEach 并使用 WriteLine 作为主体来写入值 出去。您在控制台中打印了 4, 3, 2, 1。这是因为 WriteLine 不会闲置,因此在 4 个 WriteLine 活动之后 计划时,它们使用堆栈行为(先进后出)执行。

但是,如果您的身体中有一些活动可以闲置,例如 接收活动或延迟活动。那么就无需等待 他们来完成。 ParallelForEach 转到下一个预定的主体 活动并尝试执行它。如果该活动也空闲, ParallelForEach 再次移动下一个身体活动。

我现在在哪里:

当使用 ParallelForEach 工作流活动运行我上面的“想法”时,我实现的运行时间与正常的 ForEach 工作流活动大致相同。我正在考虑使底层 BeginWorkflow 方法异步,但我不确定这对于 Windows WF 的运行方式是否是一个好主意。

我需要你的帮助:

有人对我如何实现我想要达到的结果有任何建议吗?是否有另一种方法可以在尽可能多的线程上并行执行 foreach 工作流的主体?我有 8 个逻辑处理器,我想利用它们,感觉集合的每次迭代都是独立于其他的。

有什么想法吗??

【问题讨论】:

  • 如果大部分时间是因为 I/O,那么让它并行不会给你带来这么大的改进......
  • 是的,但我已经对这些方法的实现进行了编码以使用 Async/Await,以便释放线程。但是,由于它们都在同一个线程上运行,这似乎没有帮助,因为即使线程是空闲的,主体的下一次迭代也不会执行。我相信这更多地与工作流调度器的实现方式有关。
  • 1.5 x 15000 = 6 小时 15 分钟。还是 1.5 不是正确的平均值?
  • 另外:跨工作项重用服务代理可能是一种选择。实例化代理和通道有时会产生很多开销
  • 不是正确的平均值。根据每次迭代的内容,时间可以在 5 到 10 小时之间。我正在为最坏的情况做准备。

标签: c# concurrency workflow-foundation-4


【解决方案1】:

工作流运行时是单线程的。要真正进行并行工作,您必须(以某种方式)管理自己的线程。我的猜测是你的活动只是在 Execute 方法中做他们的事情,而运行时一次只允许一个 Execute。

这是 NonblockingNativeActivity 类的代码。它对我们很有用,我希望它对你也有帮助。使用它作为您的活动的基类,而不是覆盖 Execute,覆盖 ExecuteNonblocking。如果您需要使用 Workflow 运行时,您还可以覆盖 PrepareToExecute 和 AfterExecute,但它们将是单线程的。

using System.Text;
using System.Activities.Hosting;
using System.Activities;
using System.Diagnostics;
using System.Threading.Tasks;
using System.Threading;

namespace Sample.Activities
{
    /// <summary>
    /// Class Non-Blocking Native Activity
    /// </summary>
    public abstract class NonblockingNativeActivity : NativeActivity
    {
        private Variable<NoPersistHandle> NoPersistHandle { get; set; }
        private Variable<Bookmark> Bookmark { get; set; }

        private Task m_Task;
        private Bookmark m_Bookmark;
        private BookmarkResumptionHelper m_BookmarkResumptionHelper;

        /// <summary>
        /// Allows the activity to induce idle. 
        /// </summary>
        protected override bool CanInduceIdle
        {
            get
            {
                return true;
            }
        }

        /// <summary>
        /// Prepars for Execution
        /// </summary>
        /// <param name="context"></param>
        protected virtual void PrepareToExecute(
            NativeActivityContext context)
        {
        }

        /// <summary>
        /// Executes a Non-blocking Activity
        /// </summary>
        protected abstract void ExecuteNonblocking();

        /// <summary>
        /// After Execution Completes
        /// </summary>
        /// <param name="context"></param>
        protected virtual void AfterExecute(
            NativeActivityContext context)
        {
        }

        /// <summary>
        /// Executes the Activity
        /// </summary>
        /// <param name="context"></param>
        protected override void Execute(NativeActivityContext context)
        {

            //
            //  We must enter a NoPersist zone because it looks like we're idle while our
            //  Task is executing but, we aren't really
            //
            NoPersistHandle noPersistHandle = NoPersistHandle.Get(context);
            noPersistHandle.Enter(context);

            //
            //  Set a bookmark that we will resume when our Task is done
            //
            m_Bookmark = context.CreateBookmark(BookmarkResumptionCallback);
            this.Bookmark.Set(context, m_Bookmark);
            m_BookmarkResumptionHelper = context.GetExtension<BookmarkResumptionHelper>();

            //
            //  Prepare to execute
            //
            PrepareToExecute(context);

            //
            //  Start a Task to do the actual execution of our activity
            //
            CancellationTokenSource tokenSource = new CancellationTokenSource();
            m_Task = Task.Factory.StartNew(ExecuteNonblocking, tokenSource.Token);
            m_Task.ContinueWith(TaskCompletionCallback);
        }

        private void TaskCompletionCallback(Task task)
        {
            if (!task.IsCompleted)
            {
                task.Wait();
            }

            //
            //  Resume the bookmark
            //
            m_BookmarkResumptionHelper.ResumeBookmark(m_Bookmark, null);
        }


        private void BookmarkResumptionCallback(NativeActivityContext context, Bookmark bookmark, object value)
        {
            var noPersistHandle = NoPersistHandle.Get(context);

            if (m_Task.IsFaulted)
            {
                //
                //  The task had a problem
                //
                Console.WriteLine("Exception from ExecuteNonBlocking task:");
                Exception ex = m_Task.Exception;
                while (ex != null)
                {
                    Console.WriteLine(ex.Message);
                    ex = ex.InnerException;
                }

                //
                // If there was an exception exit the no persist handle and rethrow.
                //
                if (m_Task.Exception != null)
                {
                    noPersistHandle.Exit(context);
                    throw m_Task.Exception;
                }
            }

            AfterExecute(context);

            noPersistHandle.Exit(context);
        }

        //
        //  TODO: How do we want to handle cancelations?  We can pass a CancellationToekn to the task
        //  so that we cancel the task but, maybe we can do better than that?
        //
        /// <summary>
        /// Abort Activity
        /// </summary>
        /// <param name="context"></param>
        protected override void Abort(NativeActivityAbortContext context)
        {
            base.Abort(context);
        }

        /// <summary>
        /// Cancels the Activity
        /// </summary>
        /// <param name="context"></param>
        protected override void Cancel(NativeActivityContext context)
        {
            base.Cancel(context);
        }

        /// <summary>
        /// Registers Activity Metadata
        /// </summary>
        /// <param name="metadata"></param>
        protected override void CacheMetadata(NativeActivityMetadata metadata)
        {
            base.CacheMetadata(metadata);
            this.NoPersistHandle = new Variable<NoPersistHandle>();
            this.Bookmark = new Variable<Bookmark>();
            metadata.AddImplementationVariable(this.NoPersistHandle);
            metadata.AddImplementationVariable(this.Bookmark);
            metadata.RequireExtension<BookmarkResumptionHelper>();
            metadata.AddDefaultExtensionProvider<BookmarkResumptionHelper>(() => new BookmarkResumptionHelper());
        }
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-07-28
    • 2021-07-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多