如何用PowerShell读取小型Azure Event Hub实例的全部事件(偏移量0起)
读取Azure Event Hubs所有事件的PowerShell方法
方法一:使用Az.EventHubs模块
借助Azure官方PowerShell模块是最简便的实现方式:
- 安装并导入模块(未安装时执行):
Install-Module -Name Az.EventHubs -Force -AllowClobber Import-Module Az.EventHubs
- 登录Azure账号:
Connect-AzAccount
- 获取事件中心的权限连接字符串:
$eventHubNamespace = "你的命名空间名称" $eventHubName = "你的事件中心名称" $resourceGroupName = "资源组名称" $connectionString = (Get-AzEventHubAuthorizationRule -ResourceGroupName $resourceGroupName -NamespaceName $eventHubNamespace -Name "RootManageSharedAccessKey").PrimaryConnectionString
- 从偏移量0开始读取所有事件:
# 加载依赖的.NET程序集 Add-Type -Path "C:\Program Files\WindowsPowerShell\Modules\Az.EventHubs\*\Microsoft.Azure.EventHubs.dll" Add-Type -Path "C:\Program Files\WindowsPowerShell\Modules\Az.EventHubs\*\Microsoft.Azure.ServiceBus.dll" # 创建EventHub客户端 $client = [Microsoft.Azure.EventHubs.EventHubClient]::CreateFromConnectionString($connectionString, $eventHubName) try { # 获取事件中心所有分区ID $partitionIds = $client.GetPartitionIdsAsync().GetAwaiter().GetResult() foreach ($partitionId in $partitionIds) { # 创建从偏移量0开始的接收器 $receiver = $client.CreateReceiverAsync("$DefaultConsumerGroup", $partitionId, [Microsoft.Azure.EventHubs.EventPosition]::FromStart()).GetAwaiter().GetResult() Write-Host "读取分区 $partitionId 的事件..." while ($true) { # 批量接收事件,单次最多100条,超时30秒 $events = $receiver.ReceiveAsync(100, [TimeSpan]::FromSeconds(30)).GetAwaiter().GetResult() if ($events.Count -eq 0) { Write-Host "分区 $partitionId 已无更多事件" break } foreach ($event in $events) { # 解析并输出事件内容 $body = [System.Text.Encoding]::UTF8.GetString($event.Body.Array, $event.Body.Offset, $event.Body.Count) Write-Host "事件ID: $($event.SystemProperties.EventId), 内容: $body" } } $receiver.CloseAsync().GetAwaiter().GetResult() } } finally { $client.CloseAsync().GetAwaiter().GetResult() }
方法二:直接调用REST API
如果习惯用REST请求(和你发送消息的方式一致),可通过生成SAS令牌后调用API读取:
- 生成SAS令牌的PowerShell函数:
function New-SASToken { param( [string]$Namespace, [string]$EventHubName, [string]$AccessKeyName, [string]$AccessKey ) $resourceUri = [Uri]::EscapeDataString("$Namespace.servicebus.windows.net/$EventHubName") $expires = [DateTimeOffset]::Now.AddHours(1).ToUnixTimeSeconds() $stringToSign = [Uri]::EscapeDataString($resourceUri) + "`n" + $expires $hmac = New-Object System.Security.Cryptography.HMACSHA256 $hmac.Key = [Convert]::FromBase64String($AccessKey) $signature = [Convert]::ToBase64String($hmac.ComputeHash([Text.Encoding]::UTF8.GetBytes($stringToSign))) return "SharedAccessSignature sr=$resourceUri&sig=$([Uri]::EscapeDataString($signature))&se=$expires&skn=$AccessKeyName" }
- 生成令牌并读取偏移量0的事件:
$eventHubNamespace = "你的命名空间名称" $eventHubName = "你的事件中心名称" $accessKeyName = "RootManageSharedAccessKey" $accessKey = "你的共享访问密钥" # 生成SAS令牌 $sasToken = New-SASToken -Namespace $eventHubNamespace -EventHubName $eventHubName -AccessKeyName $accessKeyName -AccessKey $accessKey # 获取所有分区ID $partitionIdsUri = "https://$eventHubNamespace.servicebus.windows.net/$eventHubName/partitions?api-version=2021-11-01" $partitionIds = (Invoke-RestMethod -Uri $partitionIdsUri -Headers @{Authorization = $sasToken} -Method Get).partition foreach ($partitionId in $partitionIds) { # 初始化读取地址(从偏移量0开始) $readUri = "https://$eventHubNamespace.servicebus.windows.net/$eventHubName/partitions/$partitionId/messages/head?api-version=2021-11-01&timeout=60&startOffset=0" Write-Host "读取分区 $partitionId 的事件..." while ($true) { try { $response = Invoke-RestMethod -Uri $readUri -Headers @{Authorization = $sasToken} -Method Get -ContentType "application/json" if ($null -eq $response -or $response.Count -eq 0) { Write-Host "分区 $partitionId 已无更多事件" break } foreach ($event in $response) { Write-Host "事件内容: $($event.Content)" } # 更新读取起始位置为最后一个事件的偏移量 $lastOffset = $response[-1].SystemProperties.Offset $readUri = "https://$eventHubNamespace.servicebus.windows.net/$eventHubName/partitions/$partitionId/messages/head?api-version=2021-11-01&timeout=60&startOffset=$lastOffset" } catch { if ($_.Exception.Response.StatusCode -eq [System.Net.HttpStatusCode]::NotFound) { Write-Host "分区 $partitionId 已无更多事件" break } throw } } }
注意事项
- 确保你的共享访问密钥拥有读取权限(默认的RootManageSharedAccessKey具备全权限)
- 可根据实际事件数量调整批量读取的参数(如单次接收条数、超时时间)
- 使用模块方法时,需保证Az.EventHubs模块版本兼容,避免程序集加载错误
内容的提问来源于stack exchange,提问作者Marc
相关产品推荐
相关产品推荐

