IEventProcessor 未从事件中心读取
IEventProcessor not reading from Event Hub
我目前正在使用 EventProcessorHost 和简单的 IEventProcessor 实现来实现事件中心 reader。我已经确认遥测数据正在使用 Paolo Salvatori 的出色 Service Bus Explorer 写入事件中心。我已成功将 EventProcessorHost 配置为将存储帐户用于租约和检查点。我可以在存储帐户中看到事件中心数据文件。我此时看到的问题是 IEventProcessor 实现似乎没有从事件中心读取任何内容。
我没有收到任何例外情况。测试控制台应用程序可以毫无问题地连接到存储帐户。我注意到我添加到构造函数中的日志记录语句从未被调用,因此看起来接收器实际上从未被创建过。我觉得我缺少一些简单的东西。谁能帮我确定我错过了什么?谢谢!
IEventProcessor 实现:
namespace Receiver
{
internal class SimpleEventProcessor : IEventProcessor
{
private Stopwatch _checkPointStopwatch;
public SimpleEventProcessor()
{
Console.WriteLine("SimpleEventProcessor created");
}
#region Implementation of IEventProcessor
public Task OpenAsync(PartitionContext context)
{
Console.WriteLine("SimpleEventProcessor initialized. Partition: '{0}', Offset: '{1}",
context.Lease.PartitionId, context.Lease.Offset);
_checkPointStopwatch = new Stopwatch();
_checkPointStopwatch.Start();
return Task.FromResult<object>(null);
}
public async Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
{
foreach (var data in messages.Select(eventData => Encoding.UTF8.GetString(eventData.GetBytes())))
{
Console.WriteLine("Message received. Partition: '{0}', Data: '{1}'", context.Lease.PartitionId,
data);
}
if (_checkPointStopwatch.Elapsed > TimeSpan.FromSeconds(30))
{
await context.CheckpointAsync();
_checkPointStopwatch.Restart();
}
}
public async Task CloseAsync(PartitionContext context, CloseReason reason)
{
Console.WriteLine("Processor shutting down. Partition '{0}', Reason: {1}", context.Lease.PartitionId,
reason);
if (reason == CloseReason.Shutdown)
{
await context.CheckpointAsync();
}
}
#endregion
}
}
测试控制台代码:
namespace EventHubTestConsole
{
internal class Program
{
private static void Main(string[] args)
{
AsyncPump.Run((Func<Task>) MainAsync);
}
private static async Task MainAsync()
{
const string eventHubConnectionString =
"Endpoint=<EH endpoint>;SharedAccessKeyName=<key name>;SharedAccessKey=<key>";
const string eventHubName = "<event hub name>";
const string storageAccountName = "<storage account name>";
const string storageAccountKey = "<valid storage key>";
var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName={0};AccountKey={1}",
storageAccountName, storageAccountKey);
Console.WriteLine("Connecting to storage account with ConnectionString: {0}", storageConnectionString);
var eventProcessorHostName = Guid.NewGuid().ToString();
var eventProcessorHost = new EventProcessorHost(
eventProcessorHostName,
eventHubName,
EventHubConsumerGroup.DefaultGroupName,
eventHubConnectionString,
storageConnectionString);
var epo = new EventProcessorOptions
{
MaxBatchSize = 100,
PrefetchCount = 1,
ReceiveTimeOut = TimeSpan.FromSeconds(20),
InitialOffsetProvider = (name) => DateTime.Now.AddDays(-7)
};
epo.ExceptionReceived += OnExceptionReceived;
await eventProcessorHost.RegisterEventProcessorAsync<SimpleEventProcessor>(epo);
Console.WriteLine("Receiving. Please enter to stop worker.");
Console.ReadLine();
}
public static void OnExceptionReceived(object sender, ExceptionReceivedEventArgs args)
{
Console.WriteLine("Event Hub exception received: {0}", args.Exception.Message);
}
}
看来问题出在您 EventProcessorOptions.PrefetchCount.
的值上
我稍微更改了您的代码,如下所示(删除 AsyncPump 并干净地关闭接收器)。我发现如果 PrefetchCount 小于 10,RegisterEventProcessorAsync 会抛出异常。
namespace EventHubTestConsole
{
internal class Program
{
private static void Main(string[] args)
{
const string eventHubConnectionString =
"Endpoint=<EH endpoint>;SharedAccessKeyName=<key name>;SharedAccessKey=<key>";
const string eventHubName = "<event hub name>";
const string storageAccountName = "<storage account name>";
const string storageAccountKey = "<valid storage key>";
var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName={0};AccountKey={1}",
storageAccountName, storageAccountKey);
Console.WriteLine("Connecting to storage account with ConnectionString: {0}", storageConnectionString);
var eventProcessorHostName = Guid.NewGuid().ToString();
var eventProcessorHost = new EventProcessorHost(
eventProcessorHostName,
eventHubName,
EventHubConsumerGroup.DefaultGroupName,
eventHubConnectionString,
storageConnectionString);
var epo = new EventProcessorOptions
{
MaxBatchSize = 100,
PrefetchCount = 10,
ReceiveTimeOut = TimeSpan.FromSeconds(20),
InitialOffsetProvider = (name) => DateTime.Now.AddDays(-7)
};
epo.ExceptionReceived += OnExceptionReceived;
eventProcessorHost.RegisterEventProcessorAsync<SimpleEventProcessor>(epo).Wait();
Console.WriteLine("Receiving. Please enter to stop worker.");
Console.ReadLine();
eventProcessorHost.UnregisterEventProcessorAsync().Wait();
}
public static void OnExceptionReceived(object sender, ExceptionReceivedEventArgs args)
{
Console.WriteLine("Event Hub exception received: {0}", args.Exception.Message);
}
}
}
我目前正在使用 EventProcessorHost 和简单的 IEventProcessor 实现来实现事件中心 reader。我已经确认遥测数据正在使用 Paolo Salvatori 的出色 Service Bus Explorer 写入事件中心。我已成功将 EventProcessorHost 配置为将存储帐户用于租约和检查点。我可以在存储帐户中看到事件中心数据文件。我此时看到的问题是 IEventProcessor 实现似乎没有从事件中心读取任何内容。
我没有收到任何例外情况。测试控制台应用程序可以毫无问题地连接到存储帐户。我注意到我添加到构造函数中的日志记录语句从未被调用,因此看起来接收器实际上从未被创建过。我觉得我缺少一些简单的东西。谁能帮我确定我错过了什么?谢谢!
IEventProcessor 实现:
namespace Receiver
{
internal class SimpleEventProcessor : IEventProcessor
{
private Stopwatch _checkPointStopwatch;
public SimpleEventProcessor()
{
Console.WriteLine("SimpleEventProcessor created");
}
#region Implementation of IEventProcessor
public Task OpenAsync(PartitionContext context)
{
Console.WriteLine("SimpleEventProcessor initialized. Partition: '{0}', Offset: '{1}",
context.Lease.PartitionId, context.Lease.Offset);
_checkPointStopwatch = new Stopwatch();
_checkPointStopwatch.Start();
return Task.FromResult<object>(null);
}
public async Task ProcessEventsAsync(PartitionContext context, IEnumerable<EventData> messages)
{
foreach (var data in messages.Select(eventData => Encoding.UTF8.GetString(eventData.GetBytes())))
{
Console.WriteLine("Message received. Partition: '{0}', Data: '{1}'", context.Lease.PartitionId,
data);
}
if (_checkPointStopwatch.Elapsed > TimeSpan.FromSeconds(30))
{
await context.CheckpointAsync();
_checkPointStopwatch.Restart();
}
}
public async Task CloseAsync(PartitionContext context, CloseReason reason)
{
Console.WriteLine("Processor shutting down. Partition '{0}', Reason: {1}", context.Lease.PartitionId,
reason);
if (reason == CloseReason.Shutdown)
{
await context.CheckpointAsync();
}
}
#endregion
}
}
测试控制台代码:
namespace EventHubTestConsole
{
internal class Program
{
private static void Main(string[] args)
{
AsyncPump.Run((Func<Task>) MainAsync);
}
private static async Task MainAsync()
{
const string eventHubConnectionString =
"Endpoint=<EH endpoint>;SharedAccessKeyName=<key name>;SharedAccessKey=<key>";
const string eventHubName = "<event hub name>";
const string storageAccountName = "<storage account name>";
const string storageAccountKey = "<valid storage key>";
var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName={0};AccountKey={1}",
storageAccountName, storageAccountKey);
Console.WriteLine("Connecting to storage account with ConnectionString: {0}", storageConnectionString);
var eventProcessorHostName = Guid.NewGuid().ToString();
var eventProcessorHost = new EventProcessorHost(
eventProcessorHostName,
eventHubName,
EventHubConsumerGroup.DefaultGroupName,
eventHubConnectionString,
storageConnectionString);
var epo = new EventProcessorOptions
{
MaxBatchSize = 100,
PrefetchCount = 1,
ReceiveTimeOut = TimeSpan.FromSeconds(20),
InitialOffsetProvider = (name) => DateTime.Now.AddDays(-7)
};
epo.ExceptionReceived += OnExceptionReceived;
await eventProcessorHost.RegisterEventProcessorAsync<SimpleEventProcessor>(epo);
Console.WriteLine("Receiving. Please enter to stop worker.");
Console.ReadLine();
}
public static void OnExceptionReceived(object sender, ExceptionReceivedEventArgs args)
{
Console.WriteLine("Event Hub exception received: {0}", args.Exception.Message);
}
}
看来问题出在您 EventProcessorOptions.PrefetchCount.
的值上我稍微更改了您的代码,如下所示(删除 AsyncPump 并干净地关闭接收器)。我发现如果 PrefetchCount 小于 10,RegisterEventProcessorAsync 会抛出异常。
namespace EventHubTestConsole
{
internal class Program
{
private static void Main(string[] args)
{
const string eventHubConnectionString =
"Endpoint=<EH endpoint>;SharedAccessKeyName=<key name>;SharedAccessKey=<key>";
const string eventHubName = "<event hub name>";
const string storageAccountName = "<storage account name>";
const string storageAccountKey = "<valid storage key>";
var storageConnectionString = string.Format("DefaultEndpointsProtocol=https;AccountName={0};AccountKey={1}",
storageAccountName, storageAccountKey);
Console.WriteLine("Connecting to storage account with ConnectionString: {0}", storageConnectionString);
var eventProcessorHostName = Guid.NewGuid().ToString();
var eventProcessorHost = new EventProcessorHost(
eventProcessorHostName,
eventHubName,
EventHubConsumerGroup.DefaultGroupName,
eventHubConnectionString,
storageConnectionString);
var epo = new EventProcessorOptions
{
MaxBatchSize = 100,
PrefetchCount = 10,
ReceiveTimeOut = TimeSpan.FromSeconds(20),
InitialOffsetProvider = (name) => DateTime.Now.AddDays(-7)
};
epo.ExceptionReceived += OnExceptionReceived;
eventProcessorHost.RegisterEventProcessorAsync<SimpleEventProcessor>(epo).Wait();
Console.WriteLine("Receiving. Please enter to stop worker.");
Console.ReadLine();
eventProcessorHost.UnregisterEventProcessorAsync().Wait();
}
public static void OnExceptionReceived(object sender, ExceptionReceivedEventArgs args)
{
Console.WriteLine("Event Hub exception received: {0}", args.Exception.Message);
}
}
}