本地Dapr PubSub(Azure Event Hubs)订阅者无法接收消息求助
Dapr PubSub基于Azure Event Hubs订阅无消息接收排查
问题背景
本地基于Azure Event Hubs实现Dapr PubSub,已完成两个项目的Dapr配置。发布者调用daprclient.PublishEventAsync方法指定pubsub名称和topic后执行成功,但订阅者应用订阅同一topic无法接收消息。
发布者代码
[HttpPost] public async Task<IActionResult> WeatherChange([FromBody] WeatherChange weatherChange) { var source = new CancellationTokenSource(); var cancellationToken = source.Token; _logger.LogInformation("Customer Order received: {@WeatherChange}", weatherChange); try { await _daprClient.PublishEventAsync(Constants.PubSubName, "myeventhub", weatherChange, cancellationToken); _logger.LogInformation("Published message: {weatherChange} ", weatherChange); } catch (Exception e) { return BadRequest("Please try again"); } return Ok("Weather Change Published"); }
订阅者代码
[Topic(Constants.PubSubName, "myeventhub")] [HttpPost("/weather")] public void PostWeathers(WeatherChange weather) { _logger.LogInformation("Weather posted in subscribed............... {WeatherChange}", weather); }
订阅者Startup配置
public void ConfigureServices(IServiceCollection services) { services.AddControllers().AddDapr(); services.AddControllers(); services.AddSwaggerGen(c => { c.SwaggerDoc("v1", new OpenApiInfo { Title = "Dapr.SubscriberService", Version = "v1" }); }); } // This method gets called by the runtime. Use this method to configure the HTTP request pipeline. public void Configure(IApplicationBuilder app, IWebHostEnvironment env) { if (env.IsDevelopment()) { app.UseDeveloperExceptionPage(); app.UseSwagger(); app.UseSwaggerUI(c => c.SwaggerEndpoint("/swagger/v1/swagger.json", "Dapr.SubscriberService v1")); } app.UseRouting(); app.UseAuthorization(); // Use Cloud Events app.UseCloudEvents(); app.UseEndpoints(endpoints => { // Map subscriber handler endpoints.MapSubscribeHandler(); endpoints.MapControllers(); }); }
PubSub组件配置
apiVersion: dapr.io/v1alpha1 kind: Component metadata: name: pubsub spec: type: pubsub.azure.eventhubs version: v1 metadata: # Either connectionString or eventHubNamespace is required # Use connectionString when *not* using Azure AD - name: connectionString value: "<Removed>" # Use eventHubNamespace when using Azure AD - name: enableEntityManagement value: "false" # The following four properties are needed only if enableEntityManagement is set to true - name: resourceGroupName value: "<Removed>" - name: subscriptionID value: "<Removed>" - name: partitionCount value: "1" - name: messageRetentionInDays value: "3" # Checkpoint store attributes - name: storageAccountName value: "<Removed>" - name: storageAccountKey value: "<Removed>" - name: storageContainerName value: "<Removed>" # Alternative to passing storageAccountKey - name: storageConnectionString value: "<Removed>" scopes: - publisher-service - subscriber-service
订阅者启动命令
dapr run --app-id subscriber-service --resources-path components --app-port 5200 --dapr-grpc-port 5201 --dapr-http-port 5280
排查建议
- 确认Event Hub实体存在:因
enableEntityManagement设为false,需手动在Azure上创建名为myeventhub的Event Hub实体,Dapr不会自动生成。 - 验证存储配置有效性:检查存储账户的连接字符串/密钥是否正确,Dapr依赖Azure存储做消费 checkpoint,配置错误会导致订阅者无法追踪消息位置。
- 检查订阅者Sidecar与应用连通性:
- 访问订阅者Dapr元数据端点
http://localhost:5280/v1.0/metadata,查看subscriptions节点下是否存在myeventhub的订阅记录。 - 手动调用订阅者
/weather接口,验证应用本身能否正常接收WeatherChange格式的数据。
- 访问订阅者Dapr元数据端点
- 核对pubsub名称一致性:确认
Constants.PubSubName的值为pubsub(与组件metadata的name字段一致),名称不匹配会导致消息发送到错误的pubsub组件。 - 查看Sidecar日志:
- 启动订阅者时添加
--log-level debug参数,查看日志中是否存在Event Hub连接或消息处理的错误信息。 - 检查发布者Sidecar日志,确认消息确实发送至目标Event Hub。
- 启动订阅者时添加
- 确认发布者App ID匹配:组件
scopes已包含publisher-service,需确保发布者启动时的--app-id参数值为publisher-service。
内容的提问来源于Stack Exchange,提问作者Vishal Shah
相关产品推荐
相关产品推荐

