【发布时间】: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