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

本地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格式的数据。
  • 核对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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 22:20:29