将C# Kafka Producer转为Windows Service时的异步调用问题排查
问题描述
我原本用Confluent库开发了一个Kafka Producer控制台应用,作为定时任务定期向Kafka Broker提交错误记录,运行正常。现在要把它转为Windows Service,处理异步await调用时遇到问题。相关代码如下:
Program.cs代码
static void Main() { /* Uncomment for debugging */ //Debugger.Launch(); ServiceBase[] ServicesToRun; ServicesToRun = new ServiceBase[] { new Kafka_Producer() }; ServiceBase.Run(ServicesToRun); }
服务API文件代码
public Kafka_Producer() { InitializeComponent(); } protected override void OnStart(string[] args) { // Some initialization, setting global variables ... producer = new Producer(sBootstrapServer); // 该行提示:"Because this call is not awaited, execution of the current method continues before the call is completed. Consider applying the await operator to the result of the call" CheckDB(); // Set up a timer that triggers periodically. timer = new System.Timers.Timer(); timer.Interval = iRunInterval; timer.Elapsed += new ElapsedEventHandler(this.OnTimer); timer.Start(); } protected async Task CheckDB() { List<string> lsJsonResult = ApplicationException.GetErrorsJSON(); await producer.SendErrorData(sTopicName, lsJsonResult); }
SendErrorData方法代码
public async Task SendErrorData(string topicName, List<string> jsonResult) { string logMessage = string.Empty; using (var producer = new ProducerBuilder<Null, string>(_producerConfig) .SetValueSerializer(Serializers.Utf8) .SetLogHandler((_, message) => LogWriter.LogWrite($"Facility: {message.Facility}-{message.Level} Message: {message.Message}")) .SetErrorHandler((_, e) => LogWriter.LogWrite($"Error: {e.Reason}. Is Fatal: {e.IsFatal}")) .Build()) try { foreach(string sJsonErrData in jsonResult) { dynamic results = JsonConvert.DeserializeObject<dynamic>(sJsonErrData); int id = results.uniqueId; var deliveryReport = await producer.ProduceAsync(topicName, new Message<Null, string> { Value = sJsonErrData }); .... } } catch() {} }
问题根源与修复方案
问题根源
- OnStart无法直接异步等待:Windows Service的
OnStart是同步方法,直接调用CheckDB()不await属于"火与遗忘"操作,会导致异步任务脱离上下文,可能出现任务未完成就被终止、异常无法捕获的情况,编译器警告正是这个原因。 - Timer事件的异步处理隐患:如果后续
OnTimer要调用异步方法,同样会出现未等待的问题,且空catch块会隐藏所有异常,无法排查问题。
修复步骤
1. 处理OnStart中的CheckDB调用
因为OnStart不能标记为async,可选择两种方式:
同步等待(适合短耗时初始化):
// 替换原CheckDB(); CheckDB().GetAwaiter().GetResult();注意:如果操作耗时过长,可能导致服务启动超时被系统终止,耗时久的操作建议用后台线程。
后台异步执行(避免阻塞启动):
// 替换原CheckDB(); _ = Task.Run(async () => { try { await CheckDB(); } catch (Exception ex) { LogWriter.LogWrite($"初始化CheckDB失败: {ex.Message}"); } });
2. 修复Timer的Elapsed事件处理
将事件处理改为async void(Timer事件是async void的合理场景),并防止定时器重入:
private async void OnTimer(object sender, ElapsedEventArgs e) { timer.Stop(); // 防止同一时间多次触发 try { await CheckDB(); } catch (Exception ex) { LogWriter.LogWrite($"定时任务执行失败: {ex.Message}"); } finally { timer.Start(); // 恢复定时器 } }
3. 完善异常捕获
把SendErrorData中的空catch块补充日志,避免隐藏异常:
catch (Exception ex) { LogWriter.LogWrite($"发送Kafka消息失败: {ex.Message}, 详情: {ex.ToString()}"); }
内容的提问来源于stack exchange,提问作者NoBullMan
相关产品推荐
相关产品推荐

