代码之家  ›  专栏  ›  技术社区  ›  Jonas Benz

如何重新启动异步方法?取消前一次运行,等待它,然后启动它

  •  0
  • Jonas Benz  · 技术社区  · 7 年前

    我有一个方法 RestartAsync 它启动了一个方法 DoSomethingAsync . 什么时候? 重新启动异步 再次调用它应该取消 剂量异步 等它完成( 剂量异步 不能同步取消,并且在上一个任务仍在进行时不应调用它)。

    我的第一个方法如下:

    public async Task RestartTest()
    {
        Task[] allTasks = { RestartAsync(), RestartAsync(), RestartAsync() } ;
        await Task.WhenAll(allTasks);
    }
    
    private async Task RestartAsync()
    {
        _cts.Cancel();
        _cts = new CancellationTokenSource();
        await _somethingIsRunningTask;
    
        _somethingIsRunningTask = DoSomethingAsync(_cts.Token);
    
        await _somethingIsRunningTask;
    }
    
    private static int _numberOfStarts;
    
    private async Task DoSomethingAsync(CancellationToken cancellationToken)
    {
        _numberOfStarts++;
        int numberOfStarts = _numberOfStarts;
    
        try
        {
            Console.WriteLine(numberOfStarts + " Start to do something...");
            await Task.Delay(TimeSpan.FromSeconds(1)); // This operation can not be cancelled.
            await Task.Delay(TimeSpan.FromSeconds(1), cancellationToken);
            Console.WriteLine(numberOfStarts + " Finished to do something...");
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine(numberOfStarts + " Cancelled to do something...");
        }
    }
    

    三次调用restartasync时的实际输出如下(请注意,第二次运行正在取消并等待第一次运行,但同时第三次运行也在等待第一次运行,而不是取消并等待第二次运行):

    1 Start to do something...
    1 Cancelled to do something...
    2 Start to do something...
    3 Start to do something...
    2 Finished to do something...
    3 Finished to do something...
    

    但我想实现的是这个输出:

    1 Start to do something...
    1 Cancelled to do something...
    2 Start to do something...
    2 Cancelled to do something...
    3 Start to do something...
    3 Finished to do something...
    

    我目前的解决方案如下:

    private async Task RestartAsync()
    {
        if (_isRestarting)
        {
            return;
        }
    
        _cts.Cancel();
        _cts = new CancellationTokenSource();
    
        _isRestarting = true;
        await _somethingIsRunningTask;
        _isRestarting = false;
    
        _somethingIsRunningTask = DoSomethingAsync(_cts.Token);
    
        await _somethingIsRunningTask;
    }
    

    然后我得到这个输出:

    1 Start to do something...
    1 Cancelled to do something...
    2 Start to do something...
    2 Finished to do something...
    

    至少现在 剂量异步 未在运行中启动(请注意,第三次运行被忽略,这并不重要,因为它应该取消第二次运行)。

    但是这个解决方案感觉不太好,我必须在我想要这种行为的任何地方重复这种丑陋的模式。 这种重启机制有什么好的模式或框架吗?

    2 回复  |  直到 7 年前
        1
  •  1
  •   Leisen Chang    7 年前

    我认为问题出在RestartAsync方法内部。请注意,如果异步方法要等待某个任务,它将立即返回该任务,因此第二个RestartAsync实际上在交换其任务之前返回,第三个RestartAsync将进入并等待任务第一个RestartAsync。

    另外,如果RestartAsync将由多个线程执行,则可能需要将CTS和SomeThingIsrunning任务包装成一个线程,并使用interlocked.exchange方法交换值,以防止出现争用情况。

    以下是我的示例代码,未完全测试:

    public class Program
    {
        static async Task Main(string[] args)
        {
            RestartTaskDemo restartTaskDemo = new RestartTaskDemo();
    
            Task[] tasks = { restartTaskDemo.RestartAsync( 1000 ), restartTaskDemo.RestartAsync( 1000 ), restartTaskDemo.RestartAsync( 1000 ) };
            await Task.WhenAll( tasks );
    
            Console.ReadLine();
        }
    }
    
    public class RestartTaskDemo
    {
        private int Counter = 0;
    
        private TaskEntry PreviousTask = new TaskEntry( Task.CompletedTask, new CancellationTokenSource() );
    
        public async Task RestartAsync( int delay )
        {            
            TaskCompletionSource<bool> taskCompletionSource = new TaskCompletionSource<bool>();
            CancellationTokenSource cancellationTokenSource = new CancellationTokenSource();
    
            TaskEntry previousTaskEntry = Interlocked.Exchange( ref PreviousTask, new TaskEntry( taskCompletionSource.Task, cancellationTokenSource ) );
    
            previousTaskEntry.CancellationTokenSource.Cancel();
            await previousTaskEntry.Task.ContinueWith( Continue );
    
            async Task Continue( Task previousTask )
            {
                try
                {
                    await DoworkAsync( delay, cancellationTokenSource.Token );
                    taskCompletionSource.TrySetResult( true );
                }
                catch( TaskCanceledException )
                {
                    taskCompletionSource.TrySetCanceled();
                }
            }            
        }
    
        private async Task DoworkAsync( int delay, CancellationToken cancellationToken )
        {
            int count = Interlocked.Increment( ref Counter );
            Console.WriteLine( $"Task {count} started." );
    
            try
            {
                await Task.Delay( delay, cancellationToken );
                Console.WriteLine( $"Task {count} finished." );
            }
            catch( TaskCanceledException )
            {
                Console.WriteLine( $"Task {count} cancelled." );
                throw;
            }
        }
    
        private class TaskEntry
        {
            public Task Task { get; }
    
            public CancellationTokenSource CancellationTokenSource { get; }
    
            public TaskEntry( Task task, CancellationTokenSource cancellationTokenSource )
            {
                Task = task;
                CancellationTokenSource = cancellationTokenSource;
            }
        }
    }
    
        2
  •  1
  •   Paulo Morgado    7 年前

    这是一个并发问题。所以,您需要一个并发问题的解决方案:信号量。

    在一般情况下,还应考虑运行的方法何时引发 OperationCanceledException :

    private async Task DoSomethingAsync(CancellationToken cancellationToken)
    {
        _numberOfStarts++;
        int numberOfStarts = _numberOfStarts;
    
        try
        {
            Console.WriteLine(numberOfStarts + " Start to do something...");
            await Task.Delay(TimeSpan.FromSeconds(1)); // This operation can not be cancelled.
            await Task.Delay(TimeSpan.FromSeconds(1), cancellationToken);
            Console.WriteLine(numberOfStarts + " Finished to do something...");
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine(numberOfStarts + " Cancelled to do something...");
            throw;
        }
    }
    

    试试这个:

    private SemaphoreSlim semaphore = new SemaphoreSlim(1);
    private (CancellationTokenSource cts, Task task)? state;
    
    private async Task RestartAsync()
    {
        Task task = null;
    
        await this.semaphore.WaitAsync();
    
        try
        {
            if (this.state.HasValue)
            {
                this.state.Value.cts.Cancel();
                this.state.Value.cts.Dispose();
    
                try
                {
                    await this.state.Value.task;
                }
                catch (OperationCanceledException)
                {
                }
    
                this.state = null;
            }
    
            var cts = new CancellationTokenSource();
            task = DoSomethingAsync(cts.Token);
    
            this.state = (cts, task);
        }
        finally
        {
            this.semaphore.Release();
        }
    
        try
        {
            await task;
        }
        catch (OperationCanceledException)
        {
        }
    }