【问题标题】:Background Threads and Tasks后台线程和任务
【发布时间】:2020-11-07 01:03:50
【问题描述】:

我正在尝试找到从专用后台线程运行任务的最佳方式。

使用上下文是从 Kafka 主题消费并引发异步事件处理程序来处理 ConsumeResult 实例。

Kafka Consumer(下面的 consumer 实例)阻塞线程,直到消息被消费或传递的 CancellationToken 被取消。

consumeThread = new Thread(Consume)
{
    Name = "Kafka Consumer Thread",
    IsBackground = true,
};

这是我想出的Consume方法的实现,由上面的专用线程启动:

private void Consume(object _)
{
    try
    {
        while (!cancellationTokenSource.IsCancellationRequested)
        {
            var consumeResult = consumer.Consume(cancellationTokenSource.Token);

            var consumeResultEventArgs = new ConsumeResultReceivedEventArgs<TKey, TValue>(
                consumer, consumeResult, cancellationTokenSource.Token);

            _ = Task.Run(async () =>
            {
                if (onConsumeResultReceived is null) continue;

                var handlerInstances = onConsumeResultReceived.GetInvocationList();
                foreach (ConsumeResultReceivedEventHandler<TKey, TValue> handlerInstance in handlerInstances)
                {
                    if (cancellationTokenSource.IsCancellationRequested) return;                        
                    await handlerInstance(this, consumeResultEventArgs).ConfigureAwait(false);                            
                }

            }, cancellationTokenSource.Token);
        }
    }
    catch (OperationCanceledException)
    {

    }
    catch (ThreadInterruptedException)
    {

    }
    catch (ThreadAbortException)
    {
        // Aborting a thread is not implemented in .NET Core.
    }
}

我不确定这是从专用线程运行任务的推荐方式,因此非常感谢任何建议。

【问题讨论】:

    标签: multithreading .net-core apache-kafka task kafka-consumer-api


    【解决方案1】:

    我不清楚为什么您需要一个专用线程。当前的代码启动一个线程,然后该线程阻塞以供使用,然后在线程池线程上引发事件处理程序。

    _ = Task.Run 成语是“一劳永逸”,从某种意义上说,它很危险,因为它会默默地吞下来自您的事件引发代码或事件处理程序的任何异常。

    我建议将Thread 替换为Task.Run,并直接引发事件处理程序:

    consumeTask = Task.Run(ConsumeAsync);
    
    private async Task ConsumeAsync()
    {
      while (true)
      {
        var consumeResult = consumer.Consume(cancellationTokenSource.Token);
        var consumeResultEventArgs = new ConsumeResultReceivedEventArgs<TKey, TValue>(
            consumer, consumeResult, cancellationTokenSource.Token);
    
        if (onConsumeResultReceived is null) continue;
    
        var handlerInstances = onConsumeResultReceived.GetInvocationList();
        foreach (ConsumeResultReceivedEventHandler<TKey, TValue> handlerInstance in handlerInstances)
        {
          if (cancellationTokenSource.IsCancellationRequested) return;
          await handlerInstance(this, consumeResultEventArgs).ConfigureAwait(false);
        }
      }
    }
    

    【讨论】:

    • 谢谢你,史蒂文。我一直非常喜欢你的博客。我一直认为要避免任务中的线程阻塞调用(例如上面的 Kafka 客户端库 Consume) - 因此是专用线程。专用线程可能会阻塞,但不会过度忙于专用于运行任务的线程池线程。但是,如果我的担忧没有根据,我很乐意使用您的解决方案。
    • 线程池将自我调整,通常在几秒钟后。因此,通过使用Task.Run,您正在从线程池中获取一个线程,并且它会在几秒钟后启动一个新线程。这对于一个非常长寿命的线程来说很好,但您不希望经常这样做 - 例如,不要在 ASP.NET 中使用 Task.Run 它将为每个请求运行。
    猜你喜欢
    • 2014-05-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-02
    • 1970-01-01
    相关资源
    最近更新 更多