【问题标题】:Is there such a synchronization tool as "single-item-sized async task buffer"?有没有“单项异步任务缓冲区”这样的同步工具?
【发布时间】:2013-12-07 03:07:01
【问题描述】:

在 UI 开发中,我多次处理事件的方式是,当事件第一次出现时 - 我会立即开始处理,但如果有一个处理操作正在进行中 - 在处理另一个事件之前,我会等待它完成。如果在操作完成之前发生了多个事件 - 我只处理最近的一个。

我通常这样做的方式是,我的 process 方法有一个循环,在我的事件处理程序中,我检查一个字段,该字段指示我当前是否正在处理某些东西,如果我正在处理 - 我将当前事件参数放在另一个字段中,该字段基本上是一个项目大小的缓冲区,当当前处理过程完成时 - 我检查是否还有其他事件要处理,然后循环直到完成。

现在这似乎有点太重复了,可能不是最优雅的方式,尽管它似乎对我来说很好用。那么我有两个问题:

  1. 我需要做的事情有名字吗?
  2. 是否有一些可重用的同步类型可以为我做到这一点?

我正在考虑在我 included in my toolkit 的 Stephen Toub 的异步协调原语集中添加一些东西。

【问题讨论】:

标签: c# asynchronous synchronization task-parallel-library async-await


【解决方案1】:

首先,我们将处理您所描述的情况,即始终从 UI 线程或其他一些同步上下文中使用该方法。 Run 方法本身可以是 async 来为我们处理通过同步上下文的所有封送处理。

如果我们正在运行,我们只需设置下一个存储的操作。如果不是,则表明我们现在正在运行,等待该动作,然后继续等待下一个动作,直到没有下一个动作。我们确保无论何时完成,我们都表明我们已完成运行:

public class EventThrottler
{
    private Func<Task> next = null;
    private bool isRunning = false;

    public async void Run(Func<Task> action)
    {
        if (isRunning)
            next = action;
        else
        {
            isRunning = true;
            try
            {
                await action();
                while (next != null)
                {
                    var nextCopy = next;
                    next = null;
                    await nextCopy();
                }
            }
            finally
            {
                isRunning = false;
            }
        }
    }

    private static Lazy<EventThrottler> defaultInstance =
        new Lazy<EventThrottler>(() => new EventThrottler());
    public static EventThrottler Default
    {
        get { return defaultInstance.Value; }
    }
}

因为这个类,至少一般来说,将专门从 UI 线程中使用,所以通常只需要一个,所以我添加了一个默认实例的便利属性,但因为它可能仍然有意义在一个程序中不止一个,我没有把它变成一个单例。

Run 接受 Func&lt;Task&gt; 并认为它通常是一个异步 lambda。它可能看起来像:

public class Foo
{
    public void SomeEventHandler(object sender, EventArgs args)
    {
        EventThrottler.Default.Run(async () =>
        {
            await Task.Delay(1000);
            //do other stuff
        });
    }
}

好的,所以,为了冗长,这里有一个版本可以处理从不同线程调用事件处理程序的情况。我知道您说过您假设它们都是从 UI 线程调用的,但我对其进行了概括。这意味着锁定对 lock 块中类型的实例字段的所有访问,但不会实际执行 lock 块内的函数。最后一部分不仅对性能很重要,以确保我们不会阻止项目仅设置next 字段,而且还避免该操作也调用运行时出现问题,因此它不需要处理重新-入口问题或潜在的死锁。这种在锁块中做事,然后根据锁中确定的条件做出响应的模式意味着设置局部变量来指示锁结束后应该做什么。

public class EventThrottlerMultiThreaded
{
    private object key = new object();
    private Func<Task> next = null;
    private bool isRunning = false;

    public void Run(Func<Task> action)
    {
        bool shouldStartRunning = false;
        lock (key)
        {
            if (isRunning)
                next = action;
            else
            {
                isRunning = true;
                shouldStartRunning = true;
            }
        }

        Action<Task> continuation = null;
        continuation = task =>
        {
            Func<Task> nextCopy = null;
            lock (key)
            {
                if (next != null)
                {
                    nextCopy = next;
                    next = null;
                }
                else
                {
                    isRunning = false;
                }
            }
            if (nextCopy != null)
                nextCopy().ContinueWith(continuation);
        };
        if (shouldStartRunning)
            action().ContinueWith(continuation);
    }
}

【讨论】:

  • 这里有两个。我认为这应该更符合您的预期。
  • 这是否完成了被忽略的任务?它似乎只是被覆盖(next = action)。是垃圾收集吗?等待它的东西会发生什么?
  • @NateDiamond 这甚至没有开始任务,因为没有调用创建它的函数。鉴于它从一开始就不存在,显然它永远不会完成,也不会被 GC 处理。也就是说,除非该函数只是指一个已经开始的任务,但是如果你已经开始了这些任务,那或多或少会破坏整个模型的目的。另请注意,作为一般规则,您实际上并不需要担心垃圾收集。项目在没有被引用时会消失,所以它几乎可以自行处理。
  • 谢谢,您拥有的单线程解决方案几乎是我通常所做的对象包装的可重用版本,并且线程安全版本看起来经过深思熟虑。当我有机会时,我需要对其进行测试,如果你同意的话,我需要将它包含在我的库中(它是 MIT 许可的)。除此之外 - 我仍然想为它找出一个比节流器更好的名称,因为它与 RX 的 Throttle 方法不同,它甚至不处理事件本身。
  • 节流是它的作用,所以也许像SingleTaskThrottleAsync?
【解决方案2】:

我需要做的事情有名字吗?

您所描述的内容听起来有点像蹦床与折叠队列的结合。蹦床基本上是一个循环,它迭代地调用 thunk-returning 函数。一个例子是反应式扩展中的CurrentThreadScheduler。当一个项目在CurrentThreadScheduler 上调度时,该工作项目被添加到调度程序的线程本地队列中,之后会发生以下情况之一:

  1. 如果蹦床已经在运行(即当前线程已经在处理线程本地队列),那么Schedule() 调用会立即返回。
  2. 如果蹦床运行(即当前线程上没有工作项排队/运行),则当前线程开始处理线程本地队列中的项目,直到它为空,此时对Schedule() 的调用返回。

折叠队列会累积要处理的项目,如果队列中已经有一个等效项目,则该项目会简单地替换为较新的项目(导致仅保留最新的等效项目)队列,而不是两者)。这个想法是避免处理陈旧/过时的事件。考虑市场数据的消费者(例如,股票报价)。如果您收到多个频繁交易证券的更新,则每次更新都会使较早的更新过时。如果最近的分时已经到达,那么处理同一证券的较早分时可能没有意义。因此,折叠队列是合适的。

在您的场景中,您基本上有一个蹦床处理一个折叠队列,所有传入事件都被认为是等效的。这导致有效的最大队列大小为 1,因为添加到非空队列中的每个项目都会导致现有项目被逐出。

是否有一些可重用的同步类型可以为我做到这一点?

我不知道现有的解决方案可以满足您的需求,但您当然可以创建一个能够支持可插入调度策略的通用蹦床或事件循环。默认策略可以使用标准队列,而其他策略可能使用优先队列或折叠队列。

【讨论】:

  • 这似乎可以通过 Reactive Extensions 进行理想的处理,尽管不幸的是我自己知道的不够多。
  • 我也有同样的想法,但我想不出一个标准的 Rx 运算符组合来提供所需的功能。我相信@PaulBetts 可以想出一个。
  • 我一直想尝试一下 RX,直到上周我花了 4 天时间尝试用它做一些简单的 2D UI 手势跟踪并最终用非 RX 做在浪费前 4 天尝试学习 RX 后,在 2 小时内编写代码。我要么太愚蠢,要么太复杂/不太适合我正在做的事情。不过,我可以向 Bart De Smet 寻求帮助... :)
【解决方案3】:

您所描述的内容听起来与 TPL Dataflow 的 BrodcastBlock 的行为方式非常相似:它始终只记住您发送给它的最后一项。如果您将它与执行您的操作并且仅对当前正在处理的项目具有容量的ActionBlock 结合使用,您将得到您想要的(该方法需要一个更好的名称):

// returns send delegate
private static Action<T> CreateProcessor<T>(Action<T> executedAction)
{
    var broadcastBlock = new BroadcastBlock<T>(null);
    var actionBlock = new ActionBlock<T>(
      executedAction, new ExecutionDataflowBlockOptions { BoundedCapacity = 1 });

    broadcastBlock.LinkTo(actionBlock);

    return item => broadcastBlock.Post(item);
}

用法可能是这样的:

var processor = CreateProcessor<int>(
    i =>
    {
        Console.WriteLine(i);
        Thread.Sleep(i);
    });

processor(100);
processor(1);
processor(2);

输出:

100
2

【讨论】:

  • 看起来很酷且简单的解决方案。如果只是为了避免对您需要安装的组件的依赖,我宁愿使用@Servy 的答案。虽然它看起来确实很有趣,可能让我也去研究那个 Dataflow 库......
  • 这也可以与任务而不是动作一起使用,这样我就可以等待任务而不是运行动作吗?
  • @FilipSkakun 是的,数据流,支持async 代表。您可以执行new ActionBlock&lt;Whatever&gt;(async x =&gt; await someAsyncMethod(x), …) 之类的操作。如果您想使用我上面的代码执行此操作,请将 exectedAction 的类型更改为 Func&lt;T, Task&gt;
猜你喜欢
  • 2021-09-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-06-08
  • 2019-03-18
  • 2017-04-20
相关资源
最近更新 更多