【发布时间】: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