使用NServiceBus SqlTransport:消息存储验证失败及队列配置咨询
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:
You used
Publishinstead ofSendPublishis designed for the publish-subscribe (pub/sub) pattern. When you publish a message, NServiceBus first checks thedbo.Subscriptionstable you configured to find endpoints that have subscribed toOrderSubmitted. 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.The
sender.Samples.Sql.Sendertable is your sender's input queue
That table exists to hold messages sent to theSamples.Sql.Senderendpoint, not messages it publishes. Since you're not sending messages to your sender endpoint itself, it stays empty.How to test message persistence
To verify messages get stored in a queue when the receiver is down, switch to usingSendto 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.Workerif 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(notPublish) 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 toShared.Worker.Queueinstead 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

