【问题标题】:TPL .Net Concurrenty Issue with Azure EventHub ProducerAzure EventHub Producer 的 TPL .Net 并发问题
【发布时间】:2020-09-24 14:09:19
【问题描述】:

我正在使用 Azure 事件中心生产者客户端并从 kafka 流中读取消息,然后将其传递给反序列化/映射,然后传递给事件中心。我有一个消耗循环,它为每个消耗创建一个任务,然后有两种方法进行处理(从 kafka 滞后的角度来看,这似乎大大提高了速度。但是,事件中心让你创建一个我不创建的事件批处理一定要使用。我现在只想一次发送一条消息。为了创建一个新批次,我必须调用 Dispose()。我遇到了一个问题,即另一个函数调用当我调用 Dispose() 时,我收到一条错误消息,指出事件中心正在使用该对象。

我也尝试过使用 eventHubProducerClient.SendAsync 的重载,它允许您传入一个 IEnumerable,但我遇到了同样的问题。

所以我认为这是一个同步问题,或者我可能需要在某处加锁?

任何帮助将不胜感激。

       public void Execute()
                {
                    using (_consumer)
                    {
                        try
                        {
                            _consumer.Subscribe(_streamConsumerSettings.Topic);
                            while (true)
                            {
                                var result = _consumer.Consume(1000);
        
                                if (result == null)
                                {
                                    continue;
                                }
                                var process = Task.Factory.StartNew(() => ProcessMessage(result?.Message?.Value));
                                var send = process.ContinueWith(t => SendMessage(process.Result));                        
                            }
        
                        }
                        catch (ConsumeException e)
                        {
                            _logger.LogError(e, e.StackTrace ?? e.Message);
                            _cancelConsume = true;
                            _consumer.Close();
                            RestartConsumer();
                        }
                    }
                }
    
         public static EquipmentJson ProcessMessage(byte[] result)
         {
                    var json = _messageProcessor.DeserializeAndMap(result);
                    return json;
         }
        
         public static void SendMessage(EquipmentJson message)
         {
                    try 
                    {   
        
                        _eventHubClient.AddToBatch(message);             
                        
                    }
                    catch (Exception e)
                    {
                        _logger.LogError(e, e.StackTrace ?? e.Message);
                    }
          }
    
     
    
        public async Task AddToBatch(EquipmentJson message)
                {
                    if 
      (!string.IsNullOrEmpty(message.EquipmentLocation))
                    {
                        try
                        {
                            var batch = await _equipmentLocClient.CreateBatchAsync();
                            batch.TryAdd(new EventData(Encoding.UTF8.GetBytes(message.EquipmentLocation)));
                            await _eventHubProducerClient.SendAsync(batch);
                            batch.Dispose();
                            _logger.LogInformation($"Data sent {DateTimeOffset.UtcNow}");
                        }
                        catch (Exception e)
                        {
                            _logger.LogError(e, e.StackTrace ?? e.Message);
                        }
                    }
                }

 public class EventHubClient : IEventHubClient
    {
        private readonly ILoggerAdapter<EventHubClient> _logger;
        private readonly EventHubClientSettings _eventHubClientSettings;
        private IMapper _mapper;

        
        private static EventHubProducerClient _equipmentLocClient;


        public EventHubClient(ILoggerAdapter<EventHubClient> logger, EventHubClientSettings eventHubClientSettings, IMapper mapper)
        {
            _logger = logger;
            _eventHubClientSettings = eventHubClientSettings;
            _mapper = mapper;
            _equipmentLocClient = new EventHubProducerClient(_eventHubClientSettings.ConnectionString, _eventHubClientSettings.EquipmentLocation);

        }
    }
}

【问题讨论】:

  • 您能帮我了解您正在使用哪个事件中心客户端库并将客户端创建包含在您的 sn-p 中吗?
  • 我正在使用 MSFT 快速入门中的最新库。 Azure.Messaging.EventHubs;使用客户端创建更新了上面的帖子。我将 EventHubClient 类注册为单例,这可能是我的问题。因为我在不同的任务中多次参加该课程。
  • System.InvalidOperationException:事件批处理当前正在与事件中心服务通信;在活动操作完成之前,可能不会添加事件。在 Azure.Messaging.EventHubs.Producer.EventDataBatch.AssertNotLocked() 在 Azure.Messaging.EventHubs.Producer.EventDataBatch.TryAdd(EventData eventData)
  • 这很有趣。我没有发现任何明显的东西。您对客户端的使用和发送流程非常好。客户端可以安全地同时使用并作为长期对象使用,并且每个发送操作都是独立的。多个 Send 调用可以同时处于活动状态,无需同步。该堆栈跟踪表明在进行发送操作时某些东西正在尝试修改批处理。我无法使用您的代码的 sn-ps 在控制台应用程序中重现该行为。
  • 目前我能提供的唯一猜测是,您的流程中似乎有一些东西正在连续运行两个等待的调用,而没有屈服于 TryAdd。

标签: c# task-parallel-library azure-eventhub


【解决方案1】:

根据我对 cme​​ts 的推测,我很好奇重构为使用 async/await 而不是主循环中的显式延续可能会有所帮助。也许类似于以下 LinqPad sn-p:

async Task Main()
{
    while (true)
    {
        var message = await Task.Factory.StartNew(() => GetText());
        var events = new[] { new EventData(Encoding.UTF8.GetBytes(message)) };
        
        await Send(events).ConfigureAwait(false);
    }
}

public EventHubProducerClient client = new EventHubProducerClient("<< CONNECTION STRING >>");

public async Task Send(EventData[] events)
{
    try
    {
        await client.SendAsync(events).ConfigureAwait(false);
        "Sent".Dump();
    }
    catch (Exception ex)
    {
        ex.Dump();
    }
}

public string GetText()
{
    Thread.Sleep(250);
    return "Test";
}

如果您打算保持延续,我想知道在延续中进行轻微的结构重构是否会有所帮助,既可以推动事件的创建,也可以兑现 await 声明。也许类似于以下 LinqPad sn-p:

async Task Main()
{
    while(true)
    {
        var t = Task.Factory.StartNew(() => GetText());
        var _ = t.ContinueWith(async q =>
        {
            var events = new[] { new EventData(Encoding.UTF8.GetBytes(t.Result)) };
            await Send(events).ConfigureAwait(false);
        });
        
        await Task.Yield();
    }
}

public EventHubProducerClient client = new EventHubProducerClient("<< CONNECTION STRING >>");

public async Task Send(EventData[] events)
{
    try
    {
        await client.SendAsync(events).ConfigureAwait(false);
        "Sent".Dump();
    }
    catch (Exception ex)
    {
        ex.Dump();
    }
}

public string GetText()
{
    Thread.Sleep(250);
    return "Test";
}

【讨论】:

  • 我会看看你的建议,让你知道。谢谢!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-12-10
  • 2011-08-20
  • 1970-01-01
  • 1970-01-01
  • 2013-06-05
  • 2014-05-14
相关资源
最近更新 更多