【问题标题】:How to disable other partitions to stop receiving messages when one partition is active in Azure event hubs如何在 Azure 事件中心中的一个分区处于活动状态时禁用其他分区以停止接收消息
【发布时间】:2020-02-07 05:05:06
【问题描述】:

我有一个安装在 3 个不同服务器中的应用程序。该应用程序订阅了一个事件中心。此事件中心有 8 个分区。所以当我在所有 3 台机器上启动我的应用程序时,所有分区都会在所有 3 台机器上随机初始化。

这样说:

VM1 : 分区 0,1,2
VM2:分区 3,4
VM3 : 分区 5,6,7

所有这些分区都在不断地接收消息。这些消息需要一个接一个地处理。现在我的要求是,在机器/服务器中,我希望一次只接收一条消息(无论初始化多少分区)。 VM1、VM2、VM3也可以并行运行。

一种情况是,在一台机器上,比如 VM1,我通过分区 0 收到一条消息。现在正在处理该消息,通常需要 15 分钟。在这 15 分钟内,我不希望分区 1 或分区 2 接收任何新消息,直到前一个消息完成。一旦之前的消息处理完成,则 3 个分区中的任何一个都准备好接收新消息。一旦任何一个分区收到另一条消息,其他分区就不应该收到任何消息。

我使用的代码是这样的:

public class SimpleEventProcessor : IEventProcessor
{
    public Task CloseAsync(PartitionContext context, CloseReason reason)
    {
       Console.WriteLine($"Processor Shutting Down. Partition '{context.PartitionId}', Reason: '{reason}'.");
       return Task.CompletedTask;
    }

    public Task OpenAsync(PartitionContext context)
    {
       Console.WriteLine($"SimpleEventProcessor initialized. Partition: '{context.PartitionId}'");
       return Task.CompletedTask;
     }

    public Task ProcessErrorAsync(PartitionContext context, Exception error)
    {
       Console.WriteLine($"Error on Partition: {context.PartitionId}, Error: {error.Message}");
       return Task.CompletedTask;
    }

    public Task 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);
          Console.WriteLine($"Message received. Partition: '{context.PartitionId}', Data: '{data}'");
          DoSomethingWithMessage(); // typically takes 15-20 mins to finish this method.
       }
       return context.CheckpointAsync();
    }
} 

这可能吗?

PS:我必须使用事件中心,没有其他选择。

【问题讨论】:

    标签: c# azure azure-eventhub


    【解决方案1】:

    您可以通过在静态锁对象上互斥来实现这一点。

        public Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
        {
            lock (lockObj)
            {
                foreach (var eventData in messages)
                {
                    var data = Encoding.UTF8.GetString(eventData.Body.Array, eventData.Body.Offset, eventData.Body.Count);
                    Console.WriteLine($"Message received. Partition: '{context.PartitionId}', Data: '{data}'");
                    DoSomethingWithMessage(); // typically takes 15-20 mins to finish this method.
                }
    
                return context.CheckpointAsync();
            }
        }
    

    不要忘记将 EventProcessorOptions.MaxBatchSize 设置为 1,如下所示。

    var epo = new EventProcessorOptions
    {
        MaxBatchSize = 1
    };
    
    await eventProcessorHost.RegisterEventProcessorAsync<MyProcessorHost>(epo);
    

    带有下游代理的完整处理器代码。

    public class SampleEventProcessor : IEventProcessor
    {
        public Task OpenAsync(PartitionContext context)
        {
            Console.WriteLine($"Opened partition {context.PartitionId}");
            return Task.FromResult<object>(null);
        }
    
        public Task CloseAsync(PartitionContext context, CloseReason reason)
        {
            Console.WriteLine($"Closed partition {context.PartitionId}");
            return Task.FromResult<object>(null);
        }
    
        public async Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
        {
            foreach (var eventData in messages)
            {
                // Process the mesasage in downstream agent.
                await DownstreamAgent.ProcessEventAsync(eventData);
    
                // Checkpoint current position.
                await context.CheckpointAsync();
            }
        }
    
        public Task ProcessErrorAsync(PartitionContext context, Exception error)
        {
            Console.WriteLine($"Partition {context.PartitionId} - {error.Message}");
            return Task.CompletedTask;
        }
    }
    
    public class DownstreamAgent
    {
        const int DegreeOfParallelism = 1;
    
        static SemaphoreSlim threadSemaphore = new SemaphoreSlim(DegreeOfParallelism, DegreeOfParallelism);
    
        public static async Task ProcessEventAsync(EventData message)
        {
            // Wait for a spot so this message can get processed.
            await threadSemaphore.WaitAsync();
    
            try
            {
                // Process the message here
                var data = Encoding.UTF8.GetString(message.Body.Array);
                Console.WriteLine(data);
            }
            finally
            {
                // Release the semaphore here so that next message waiting can be processed.
                threadSemaphore.Release();
            }
        }
    }
    

    【讨论】:

    • 感谢您的回复。我将尝试测试您提供的方法。我只是想确认它不应该阻止其他虚拟机接收消息。这种阻塞应该只发生在虚拟机中。
    • 我必须在哪里将 EventProcessorOptions.MaxBatchSize 设置为 1 ?
    • 再次感谢。但是 lockObj 是什么?
    • 不确定这是否可行。您不能保证给定的 vm 将始终处理相同的分区。像这样的锁将绑定到特定的进程。如果另一个 vm 开始处理之前由另一台机器处理过的分区,例如当租约到期时,此方法将失败。
    • 另外,在这个例子中会有一个事件处理器的三个实例,所以锁必须是共享的。但是您需要为另一台机器上的进程使用分布式锁。还不如在它后面放一个队列。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-08-13
    • 1970-01-01
    • 2018-01-18
    • 2019-10-19
    • 1970-01-01
    • 2014-12-20
    • 1970-01-01
    相关资源
    最近更新 更多