如何在不修改同步方法签名时同步执行Kafka异步DeliveryHandler
问题场景
我有一个向Kafka主题生产数据的同步Send方法(签名为void,无法修改),调用Kafka的Produce方法时传入了异步的DeliveryHandler。运行单元测试时发现代码无法进入DeliveryHandler,即便尝试用DeliveryHandler(deliveryReport).GetAwaiter().GetResult()调用也无效。需要确保DeliveryHandler的逻辑能被执行,且不能修改Send方法的签名为async Task。
相关代码
Send方法代码
public void Send( ProducerMessage<TKey, TValue> producerMessage, string topic, Action<McFlowProducerResult<TKey, TValue>> callback = default) { try { var kafkaProducerMessage = new Message<string, string>(); // DeliveryHanlder逻辑未执行? _producer.Produce( topic, kafkaProducerMessage, deliveryReport => DeliveryHandler(deliveryReport)); // TODO: 如何在不使用async await Task的情况下确保DeliveryHandler逻辑执行? } catch (Exception ex) { // 异常处理逻辑 } }
DeliveryHandler代码
// TODO: 代码从未进入这个方法 private async Task DeliveryHandler(DeliveryReport<string, string> deliveryReport) { var producerResult = new ProducerResult<string, string>(deliveryReport); if (!deliveryReport.Error.IsError) { _logger.LogError("消息成功发送到DLQ主题"); return; } _logger.LogError("无法将消息发送到DLQ主题: {0}。错误原因: {1}", deliveryReport.Topic, deliveryReport.Error.Reason); if (deliveryReport.Error.Code == ErrorCode.NetworkException) { _logger.LogError("正在将消息发送到DynamoDb"); await _fatalErrorHandler.HandleError(producerResult); } }
单元测试代码
[Fact] public void ValidateDeliveryHandlerIsInvoked() { var producerMessage = new ProducerMessage<string, string>( "aKey", "aValue", new Headers(), Timestamp.Default, 0 ); ProducerResult<string, string> callbackResult = null; _mcFlowDlqProducer.Send(producerMessage, _topicName, (mcFlowProducerResult) => { callbackResult = mcFlowProducerResult; }); Assert.NotEmpty(callbackResult.Topic); }
解决方案
1. 同步等待异步DeliveryHandler完成
Kafka的Produce回调是同步委托,直接调用异步的DeliveryHandler会导致Task被fire-and-forget,测试进程可能在Task执行前就结束。需要在回调里同步等待Task完成,同时避免上下文死锁:
修改Send方法中的Produce回调:
_producer.Produce( topic, kafkaProducerMessage, deliveryReport => DeliveryHandler(deliveryReport) .ConfigureAwait(false) .GetAwaiter() .GetResult());
2. 单元测试中添加同步等待逻辑
单元测试会在Send调用后立即执行断言,此时DeliveryHandler可能还没触发。可以用ManualResetEventSlim等待回调执行完成:
修改单元测试代码:
[Fact] public void ValidateDeliveryHandlerIsInvoked() { var producerMessage = new ProducerMessage<string, string>( "aKey", "aValue", new Headers(), Timestamp.Default, 0 ); ProducerResult<string, string> callbackResult = null; var resetEvent = new ManualResetEventSlim(false); _mcFlowDlqProducer.Send(producerMessage, _topicName, (mcFlowProducerResult) => { callbackResult = mcFlowProducerResult; resetEvent.Set(); // 回调触发后通知等待线程 }); // 等待回调执行,设置超时避免无限等待 Assert.True(resetEvent.Wait(TimeSpan.FromSeconds(5)), "回调未在超时时间内触发"); Assert.NotEmpty(callbackResult.Topic); }
3. 模拟Kafka生产者(推荐用于单元测试)
如果使用真实Kafka生产者,单元测试中可能因为生产者未完成flush就结束。可以用Mock框架(如Moq)模拟生产者,让Produce方法直接同步触发DeliveryHandler:
// 模拟IProducer接口 var mockProducer = new Mock<IProducer<string, string>>(); mockProducer.Setup(p => p.Produce( It.IsAny<string>(), It.IsAny<Message<string, string>>(), It.IsAny<Action<DeliveryReport<string, string>>>())) .Callback((string topic, Message<string, string> msg, Action<DeliveryReport<string, string>> handler) => { // 直接触发回调,模拟消息发送结果 handler(new DeliveryReport<string, string> { Topic = topic, Error = ErrorCode.NoError, // 按需设置其他属性 }); }); // 将模拟的生产者注入到_mcFlowDlqProducer实例中 _mcFlowDlqProducer = new McFlowDlqProducer(mockProducer.Object, _logger.Object, _fatalErrorHandler.Object);
内容的提问来源于stack exchange,提问作者Vismark Juarez
相关产品推荐
相关产品推荐

