You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在不修改同步方法签名时同步执行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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 22:45:34