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

使用NServiceBus SqlTransport:消息存储验证失败及队列配置咨询

Answers to Your NServiceBus SqlTransport Questions

Where Did the OrderSubmitted Messages Go?

Let's start with the mystery of your missing messages—this boils down to the difference between Publish vs Send in NServiceBus, and how SqlTransport handles each:

  1. You used Publish instead of Send
    Publish is designed for the publish-subscribe (pub/sub) pattern. When you publish a message, NServiceBus first checks the dbo.Subscriptions table you configured to find endpoints that have subscribed to OrderSubmitted. Since you disabled the receiver endpoint, there are no subscriptions in that table. Without any destination endpoints, NServiceBus discards the message immediately—it doesn't store it in any queue table.

  2. The sender.Samples.Sql.Sender table is your sender's input queue
    That table exists to hold messages sent to the Samples.Sql.Sender endpoint, not messages it publishes. Since you're not sending messages to your sender endpoint itself, it stays empty.

  3. How to test message persistence
    To verify messages get stored in a queue when the receiver is down, switch to using Send to target a specific worker queue:

    // Replace Publish with Send to a worker queue
    await endpointInstance.Send("Samples.Sql.Worker", orderSubmitted)
        .ConfigureAwait(false);
    

    With this change, even if the worker is offline, the message will be stored in the corresponding queue table (e.g., dbo.Samples.Sql.Worker if the worker uses the default schema).

Does NServiceBus Support Single-Queue Multi-Consumer (Producer/Consumer) Pattern?

Absolutely! This is called competing consumers in NServiceBus, and it's perfect for your scenario where multiple workers pull messages from a single queue, with each message processed by only one worker. Here's how to configure it:

Step 1: Configure the Producer

  • Use Send (not Publish) to send messages to a shared queue (e.g., Shared.Worker.Queue).
  • Update your producer's transport configuration to reference this shared queue:
    var transport = endpointConfiguration.UseTransport<SqlServerTransport>();
    transport.ConnectionString(@"Data Source=MySqlServer;Initial Catalog=NServiceBus;Integrated Security=True;Max Pool Size=100");
    transport.DefaultSchema("sender");
    // Map the shared queue to its schema (e.g., dbo)
    transport.UseSchemaForQueue("Shared.Worker.Queue", "dbo");
    // Keep other settings like error/audit queues as-is
    

Step 2: Configure Each Worker Instance

Every worker server will use this configuration to listen to the same shared queue:

public static async Task Main()
{
    Console.Title = "Worker.Instance1"; // Unique title for each worker instance
    var endpointConfiguration = new EndpointConfiguration("Worker.Instance1"); // Unique endpoint name per instance (for logging/monitoring)
    endpointConfiguration.SendFailedMessagesTo("error");
    endpointConfiguration.EnableInstallers();

    var transport = endpointConfiguration.UseTransport<SqlServerTransport>();
    transport.ConnectionString(@"Data Source=MySqlServer;Initial Catalog=NServiceBus;Integrated Security=True;Max Pool Size=100");
    // Override the local input queue to point to the shared queue
    transport.OverrideLocalAddress("Shared.Worker.Queue");
    // Match the schema used for the shared queue
    transport.UseSchemaForQueue("Shared.Worker.Queue", "dbo");
    transport.UseSchemaForQueue("error", "dbo");
    transport.UseSchemaForQueue("audit", "dbo");

    // Disable pub/sub features since we don't need them
    endpointConfiguration.DisableFeature<SubscriptionStorage>();
    endpointConfiguration.DisableFeature<PublishSubscribe>();

    // Use SqlPersistence if you need durable worker state, or InMemoryPersistence for simplicity
    endpointConfiguration.UsePersistence<InMemoryPersistence>();

    var endpointInstance = await Endpoint.Start(endpointConfiguration)
        .ConfigureAwait(false);

    Console.WriteLine("Worker running—press any key to exit");
    Console.ReadKey();

    await endpointInstance.Stop().ConfigureAwait(false);
}

Key Notes for This Setup:

  • All workers share the same queue: By using OverrideLocalAddress, every worker listens to Shared.Worker.Queue instead of their own endpoint-named queue.
  • Automatic competing consumption: SqlTransport handles the locking logic internally—when a worker picks up a message, it locks it so no other worker can process it. If the worker fails, the lock expires and another worker can retry the message.
  • Avoid pub/sub: Disable pub/sub features to simplify configuration and prevent unintended behavior.
  • Queue creation: Enable EnableInstallers() so NServiceBus automatically creates the shared queue table if it doesn't exist, or create it manually using SqlTransport's table schema definitions.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:17:53