【问题标题】:Directing messages to consumers将消息定向到消费者
【发布时间】:2019-11-24 04:28:16
【问题描述】:

我的客户端正在尝试向接收者发送消息。但是我注意到接收器有时不会收到客户端发送的所有消息,因此丢失了一些消息(不确定问题出在哪里?客户端或接收器)。 关于为什么会发生这种情况的任何建议。这就是我目前正在做的事情

在接收方这就是我正在做的事情。

这是事件处理器

        async Task IEventProcessor.ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
        {
            foreach (var eventData in messages)
            {
                var data = Encoding.UTF8.GetString(eventData.Body.Array, eventData.Body.Offset, eventData.Body.Count);
            }
        }

这是客户端连接到事件中心的方式

var StrBuilder = new EventHubsConnectionStringBuilder(eventHubConnectionString)
{
 EntityPath = eventHubName,
};
this.eventHubClient = EventHubClient.CreateFromConnectionString(StrBuilder.ToString());

如何将我的消息发送给特定的消费者

【问题讨论】:

  • 能否提供完整的发送客户端代码?在发送客户端,如何调用自定义方法 Send(string content)?在接收端,我看到你在 CloseAsync() 方法中调用 CheckpointAsync(),但在官方文档中,CheckpointAsync() 是在 ProcessEventsAsync 中设置的() 方法。如果您没有一些自定义要求,您应该按照official sample 进行发送/接收。
  • 是的,让我更新代码
  • @IvanYang 用户调用 SendMessage,后者调用 Send。我非常怀疑在 CloseAsync() 中放错 CheckpointAsync() 是负责任的,因为在发送消息时不会调用 CloseAsync()。
  • 你能分享一下 AsyncDispatchEvent 方法吗?我刚试过,这里不能复制。
  • 你的意思是直接调用SendMessage()方法,每次只发送一条消息?如果是这样,我会按照你的代码来调试它。

标签: c# azure azure-eventhub


【解决方案1】:

我正在使用来自 eventthub 官方文档的示例代码,用于 sendingreceiving

我有 2 个消费者群体:$Defaultnewcg。假设你有 2 个客户端,client_1 使用默认消费组($Default),client_2 使用另一个消费组(newcg)

首先,在创建发送客户端后,在SendMessagesToEventHub 方法中,我们需要添加一个带值的属性。该值应该是消费者组名称。示例代码如下:

    private static async Task SendMessagesToEventHub(int numMessagesToSend)
    {
        for (var i = 0; i < numMessagesToSend; i++)
        {
            try
            {
                var message = "444 Message";
                Console.WriteLine($"Sending message: {message}");
                EventData mydata = new EventData(Encoding.UTF8.GetBytes(message));

                //here, we add a property named "cg", it's value is the consumer group. By setting this property, then we can read this message via this specified consumer group.
                mydata.Properties.Add("cg", "newcg");

                await eventHubClient.SendAsync(mydata);

            }
            catch (Exception exception)
            {
                Console.WriteLine($"{DateTime.Now} > Exception: {exception.Message}");
            }

            await Task.Delay(10);
        }

        Console.WriteLine($"{numMessagesToSend} messages sent.");
    }

然后在client_1中,创建receiver项目后,使用default consumer group($Default) -> 在SimpleEventProcessor 类-> ProcessEventsAsync 方法中,我们可以过滤掉不必要的事件数据。 ProcessEventsAsync 方法的示例代码:

        public Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
        {
            foreach (var eventData in messages)
            {
                //filter the data here
                if (eventData.Properties["cg"].ToString() == "$Default")
                {                    
                    var data = Encoding.UTF8.GetString(eventData.Body.Array, eventData.Body.Offset, eventData.Body.Count);

                    Console.WriteLine($"Message received. Partition: '{context.PartitionId}', Data: '{data}'");
                    Console.WriteLine(context.ConsumerGroupName);
                }
            }

            return context.CheckpointAsync();
        }

而在另一个客户端,比如client_2,它使用另一个消费者组,比如它的名字是newcg,我们可以按照client_1中的步骤,只是在ProcessEventsAsync方法上稍作改动,如下:

            public Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
            {
                foreach (var eventData in messages)
                {
                    //filter the data here, using another consumer group name
                    if (eventData.Properties["cg"].ToString() == "newcg")
                    {  
                       //other code
                    }
                   }

                 return context.CheckpointAsync();
               }

【讨论】:

  • 是的,看起来我需要消费者。谢谢你的建议。
【解决方案2】:

仅当有 2 个或更多事件处理器主机从同一消费者组读取时才会发生这种情况。

如果您有 32 个分区和 2 个事件处理器主机从同一消费者组读取的事件中心。然后每个事件处理器主机将从 16 个分区中读取,依此类推。

类似地,如果 4 个事件处理器主机并行读取同一消费者组,则每个将从 8 个分区读取。

检查您是否有 2 个或更多事件处理器主机在同一个使用者组上运行。

【讨论】:

    【解决方案3】:

    我已经测试了你的代码并稍微修改了它(EventProcessorHost构造函数的不同重载,并在消费消息后添加了CheckpointAsync),然后做了一些测试。

    通过使用默认实现和默认 EventProcessorOptions(EventProcessorOptions.DefaultOptions) 我可以说我在使用消息时确实遇到了一些延迟,但所有消息都已成功处理。 所以,有时我似乎没有从某个分区收到消息,但经过一段时间后,所有消息都到达

    Here 你可以找到对我有用的实际修改代码。这是一个简单的控制台应用程序,如果有东西到达,它会打印到控制台。

            string processorHostName = Guid.NewGuid().ToString();
            var Options = new EventProcessorOptions()
            {
                MaxBatchSize = 1, //not required to make it working, just for testing
            };
            Options.SetExceptionHandler((ex) =>
            {
                System.Diagnostics.Debug.WriteLine($"Exception : {ex}");
            });
            var eventHubCS = "event hub connection string";
            var storageCS = "storage connection string";
            var containerName = "test";
            var eventHubname = "test2";
            EventProcessorHost eventProcessorHost = new EventProcessorHost(eventHubname, "$Default", eventHubCS, storageCS, containerName);
            eventProcessorHost.RegisterEventProcessorAsync<MyEventProcessor>(Options).Wait();
    

    为了将消息发送到事件中心并进行测试,我使用了 message publisher app

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-06-08
      • 2019-10-23
      • 1970-01-01
      • 1970-01-01
      • 2017-09-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多