You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用C#创建Event Hub接收器获取指定时间段内的消息?

仅获取Event Hub过去1小时消息并停止监听的修改方案

你可以通过指定读取起始位置、控制读取终止时机来实现需求,以下是修改后的代码及关键说明:

using Azure.Messaging.EventHubs.Consumer;
using System.Text;
using System.Threading;

string connectionString = "";
string eventHubName = "";
string consumerGroup = EventHubConsumerClient.DefaultConsumerGroupName;

// 计算1小时前的UTC时间(Event Hub消息入队时间为UTC时区)
var startTime = DateTimeOffset.UtcNow.AddHours(-1);
// 设置读取起始位置:仅读取入队时间晚于1小时前的消息
var startingPosition = EventPosition.FromEnqueuedTime(startTime);

// 创建取消令牌,用于终止读取操作
using var cancellationTokenSource = new CancellationTokenSource();

await using (var consumerClient = new EventHubConsumerClient(consumerGroup, connectionString, eventHubName))
{
    // 获取Event Hub的所有分区ID
    var partitionIds = await consumerClient.GetPartitionIdsAsync();
    
    // 遍历每个分区读取消息
    foreach (var partitionId in partitionIds)
    {
        await foreach (var partitionEvent in consumerClient.ReadEventsFromPartitionAsync(
            partitionId, 
            startingPosition, 
            cancellationToken: cancellationTokenSource.Token))
        {
            if (partitionEvent.Data == null)
            {
                // 当前分区已无历史消息,跳出该分区的读取循环
                break;
            }
            
            // 解析并输出消息内容
            string jsonContent = Encoding.UTF8.GetString(partitionEvent.Data.Body.ToArray());
            Console.WriteLine(jsonContent);
            
            // 若读取到的消息入队时间超过当前时间(避免读取过程中新增的消息),终止所有读取
            if (partitionEvent.Data.EnqueuedTime > DateTimeOffset.UtcNow)
            {
                cancellationTokenSource.Cancel();
                break;
            }
        }
    }
    
    // 主动触发取消,确保停止所有监听
    cancellationTokenSource.Cancel();
}

关键修改点说明

  • 精准过滤消息时间范围:通过EventPosition.FromEnqueuedTime(startTime)指定读取起始点,直接跳过1小时之前的所有消息,无需手动过滤每条消息。
  • 终止监听新消息:遍历每个分区读取历史消息,当分区无更多历史消息(partitionEvent.Data == null)或读取到实时新增的消息时,触发取消令牌终止所有读取操作,避免持续监听。
  • 时区一致性:使用UTC时间计算起始点,与Event Hub存储的消息入队时间时区保持一致,避免时区偏差导致的消息漏读或多读。

内容的提问来源于stack exchange,提问作者Alan Henry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 12:22:26