【问题标题】:Parallel Handling Kafka messages in a Windows Service using C#使用 C# 在 Windows 服务中并行处理 Kafka 消息
【发布时间】:2019-06-13 13:32:01
【问题描述】:

我有用 C# 编写的 Windows 服务。早些时候,我们使用具有多个分区的事件中心来进行消息队列。我们最近搬到了卡夫卡。为了在 c# 中实现事件中心,我们有 IEventProcessor.ProcessEventsAsync ,它会持续监听事件中心通知,并在消息发布到事件中心时触发,事件中心在后台异步运行

我在 Kafka 中没有找到任何等效的方法。

我这里的要求是订阅一个 Kafka 主题并持续消费消息。当一条消息被消费时,还应该对该消息执行一些其他操作。对于每条消息说执行时间大约需要 15 分钟,我希望 Kafka 消费者消费所有消息并将其保留在队列中,就像它接收并写入文件时一样。其他进程应该读取文件,选择消息并执行其他操作。我希望所有这些都同时/并行运行。

PS:我编写了一个控制台应用程序,它可以生成和使用一条消息。我正在寻找的是排队和并行性。

【问题讨论】:

    标签: c# .net apache-kafka


    【解决方案1】:

    对于并行性,Kafka 实现了所谓的consumer groups。 Kafka 存储“偏移量”(读取:跨主题的记录键),还存储给定消费者组在处理记录时所在的偏移量。这应该允许您使用相同的程序动态创建新的消费者实例,并通过更改组允许两个程序并行使用相同的数据来执行不同的任务。

    当我创建第一个消费者时,我发现此链接也很有帮助,以防您找到一种无需 groupId 的方法来创建它:http://cloudurable.com/blog/kafka-tutorial-kafka-consumer/index.html

    希望这会有所帮助!

    【讨论】:

      【解决方案2】:

      看看 Silverback:https://silverback-messaging.net。它抽象了许多这些问题,基本用法就这么简单:

      public class Startup
      {
          public void ConfigureServices(IServiceCollection services)
          {
              services
                  .AddSilverback()
                  .WithConnectionToMessageBroker(options => options.AddKafka())
                  .AddKafkaEndpoints(
                      endpoints => endpoints
                          .Configure(
                              config =>
                              {
                                  config.BootstrapServers = "localhost:9092";
                              })
                      .AddInbound(
                          endpoint => endpoint
                              .ConsumeFrom("my-topic")
                              .DeserializeJson(serializer => serializer.UseFixedType<SomeMessage>())
                              .Configure(
                                  config =>
                                  {
                                      config.GroupId = "test-consumer-group";
                                      config.AutoOffsetReset = AutoOffsetReset.Earliest;
                                  })))
      
                  .AddSingletonSubscriber<MySubscriber>();
          }
      }
      
      public class MySubscriber
      {
          public Task OnMessageReceived(SomeMessage message)
          {
              // TODO: process message
          }
      }
      

      【讨论】:

        猜你喜欢
        • 2012-04-12
        • 2013-05-22
        • 2010-09-27
        • 2019-07-19
        • 2022-10-18
        • 2019-08-04
        • 2013-09-23
        • 2019-06-21
        • 2016-06-06
        相关资源
        最近更新 更多