NodeJS中向Service Bus Topic发消息时如何传入PartitionKey或MessageId?
问题描述
我有一个启用了**分区(Partitioning)和重复检测(Duplicate detection)**的Service Bus Topic,按照微软官方文档的方法向该Topic批量发送消息时,出现如下错误:
Error occurred: ServiceBusError: InvalidOperationError: Message to a partitioned entity with duplicate detection enabled must have either PartitionKey or MessageId
核心问题:
使用@azure/service-bus包的Node.js应用中,如何给批量消息传入PartitionKey或MessageId?
解决方案
要解决这个错误,只需给每个消息对象添加partitionKey或messageId属性即可,两种方式任选其一:
- 添加
messageId:生成唯一标识,既满足重复检测的去重逻辑,Service Bus也会自动基于messageId的哈希值分配分区 - 添加
partitionKey:指定统一或自定义的分区键,相同键的消息会被路由到同一个分区,适合需要顺序消费的场景
修改后的代码示例
可以直接在消息数组中预定义属性,也可以在批量处理前动态添加:
const { ServiceBusClient } = require("@azure/service-bus"); const connectionString = "<SERVICE BUS NAMESPACE CONNECTION STRING>" const topicName = "<TOPIC NAME>"; // 方式1:给每个消息添加唯一的messageId(推荐用UUID生成更可靠的唯一值) const messages = [ { body: "Albert Einstein", messageId: "msg-1" }, { body: "Werner Heisenberg", messageId: "msg-2" }, { body: "Marie Curie", messageId: "msg-3" }, { body: "Steven Hawking", messageId: "msg-4" }, { body: "Isaac Newton", messageId: "msg-5" }, { body: "Niels Bohr", messageId: "msg-6" }, { body: "Michael Faraday", messageId: "msg-7" }, { body: "Galileo Galilei", messageId: "msg-8" }, { body: "Johannes Kepler", messageId: "msg-9" }, { body: "Nikolaus Kopernikus", messageId: "msg-10" } ]; // 方式2:给所有消息指定统一的partitionKey // const messages = [ // { body: "Albert Einstein", partitionKey: "physicists" }, // { body: "Werner Heisenberg", partitionKey: "physicists" }, // // 其余消息同理... // ]; async function main() { const sbClient = new ServiceBusClient(connectionString); const sender = sbClient.createSender(topicName); try { let batch = await sender.createMessageBatch(); for (let i = 0; i < messages.length; i++) { if (!batch.tryAddMessage(messages[i])) { await sender.sendMessages(batch); batch = await sender.createMessageBatch(); if (!batch.tryAddMessage(messages[i])) { throw new Error("消息过大,无法加入批量"); } } } await sender.sendMessages(batch); console.log(`已向主题 ${topicName} 发送一批消息`); await sender.close(); } finally { await sbClient.close(); } } main().catch((err) => { console.log("发生错误: ", err); process.exit(1); });
内容的提问来源于stack exchange,提问作者GThree
相关产品推荐
相关产品推荐

