如何在消息发送至Azure Service Bus主题时触发API调用?
Azure Service Bus主题集成React前端与Node API实现方案
背景
我正在开发Azure Service Bus主题相关功能,已按照官方文档编写了发送端(sender)和接收端(receiver)代码,当前代码运行正常。但需要替换批量测试数据以实现实际业务场景:将发送端代码集成到React前端应用中,在表单提交完成后,将唯一表单ID发送至主题,并触发针对该表单ID的Node API调用,目前无法理清实现逻辑,求指导。
现有发送端代码
const { ServiceBusClient } = require("@azure/service-bus"); const connectionString = "<SERVICE BUS NAMESPACE CONNECTION STRING>" const topicName = "<TOPIC NAME>"; const messages = [ { body: "Albert Einstein" }, { body: "Werner Heisenberg" }, { body: "Marie Curie" }, { body: "Steven Hawking" }, { body: "Isaac Newton" }, { body: "Niels Bohr" }, { body: "Michael Faraday" }, { body: "Galileo Galilei" }, { body: "Johannes Kepler" }, { body: "Nikolaus Kopernikus" } ]; async function main() { // 创建Service Bus客户端 const sbClient = new ServiceBusClient(connectionString); // 创建主题发送器 const sender = sbClient.createSender(topicName); try { // 尝试批量发送所有消息,如果消息太大则会失败 // await sender.sendMessages(messages); // 创建消息批次对象 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("Message too big to fit in a batch"); } } } // 发送最后一个批次 await sender.sendMessages(batch); console.log(`已向主题 ${topicName} 发送一批消息`); // 关闭发送器 await sender.close(); } finally { await sbClient.close(); } } // 执行主函数 main().catch((err) => { console.log("发生错误: ", err); process.exit(1); });
现有接收端代码
const { delay, ServiceBusClient, ServiceBusMessage } = require("@azure/service-bus"); const axios = require("axios").default; const connectionString = "<ConnectionString>" const topicName = "<TopicName>"; const subscriptionName = "<Subscription>"; async function main() { // 创建Service Bus客户端 const sbClient = new ServiceBusClient(connectionString); // 创建订阅接收器 const receiver = sbClient.createReceiver(topicName, subscriptionName); // 消息处理函数 const myMessageHandler = async (messageReceived) => { console.log(`收到消息: ${messageReceived.body}`); const response = axios({ method: 'post', url: 'http://localhost:8080/gitWrite?userprojectid=63874e2e3981e40a6f4e04a7', }); console.log(response); }; // 错误处理函数 const myErrorHandler = async (error) => { console.log(error); }; // 订阅消息并指定处理函数 receiver.subscribe({ processMessage: myMessageHandler, processError: myErrorHandler }); // 等待足够时间以接收消息 await delay(5000); await receiver.close(); await sbClient.close(); } // 执行主函数 main().catch((err) => { console.log("发生错误: ", err); process.exit(1); });
实现步骤
1. 前端React集成发送逻辑
核心注意事项
绝对不能在前端暴露Service Bus连接字符串,前端属于公开环境,直接使用会导致密钥泄露,必须通过自建后端API中转发送消息。
React表单组件示例
import { useState } from 'react'; const FormComponent = () => { const [formData, setFormData] = useState({ /* 定义你的表单字段 */ }); const handleSubmit = async (e) => { e.preventDefault(); // 步骤1:提交表单到自建后端,获取唯一表单ID const submitRes = await fetch('/api/submit-form', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(formData) }); const { formId } = await submitRes.json(); // 步骤2:调用后端中转接口,发送表单ID到Service Bus await fetch('/api/send-to-servicebus', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ formId }) }); alert('表单提交完成,后续处理已触发'); }; return ( <form onSubmit={handleSubmit}> {/* 表单输入字段 */} <input type="text" value={formData.name} onChange={(e) => setFormData({...formData, name: e.target.value})} placeholder="输入名称" /> <button type="submit">提交表单</button> </form> ); }; export default FormComponent;
2. 后端中转API实现(Node.js + Express)
创建安全的后端接口,处理前端请求并调用Service Bus发送消息:
const express = require('express'); const { ServiceBusClient } = require("@azure/service-bus"); const app = express(); app.use(express.json()); // Service Bus配置(仅后端存储,绝不暴露) const connectionString = "<你的Service Bus连接字符串>"; const topicName = "<你的主题名称>"; // 接口1:处理表单提交,返回唯一表单ID app.post('/api/submit-form', async (req, res) => { // 这里实现表单数据存储逻辑,生成唯一ID const formId = `FORM-${Date.now()}-${Math.random().toString(36).slice(2, 9)}`; // 省略存储表单数据到数据库的代码 res.json({ formId }); }); // 接口2:将表单ID发送到Service Bus app.post('/api/send-to-servicebus', async (req, res) => { const { formId } = req.body; if (!formId) return res.status(400).json({ error: '缺少表单ID' }); let sbClient = null; let sender = null; try { sbClient = new ServiceBusClient(connectionString); sender = sbClient.createSender(topicName); // 发送单条表单ID消息 await sender.sendMessages({ body: formId }); res.json({ success: true, message: '表单ID已发送至Service Bus' }); } catch (err) { console.error('发送消息失败:', err); res.status(500).json({ error: '发送消息失败' }); } finally { if (sender) await sender.close(); if (sbClient) await sbClient.close(); } }); app.listen(3001, () => console.log('后端中转服务运行在端口3001'));
3. 修改接收端代码,动态处理表单ID
调整接收端的消息处理逻辑,使用收到的表单ID调用目标API,并优化服务运行模式:
const myMessageHandler = async (messageReceived) => { const formId = messageReceived.body; console.log(`收到待处理表单ID: ${formId}`); try { // 动态使用表单ID调用Node API const response = await axios({ method: 'post', url: `http://localhost:8080/gitWrite?userprojectid=${formId}`, }); console.log(`表单${formId}处理完成,API返回:`, response.data); // 手动标记消息处理完成,避免重复消费 await messageReceived.complete(); } catch (err) { console.error(`表单${formId}处理失败:`, err); // 处理失败时将消息移入死信队列,便于后续排查 await messageReceived.deadLetter(); } }; // 移除原代码中的延迟关闭逻辑,保持服务持续监听 // await delay(5000); // await receiver.close(); // await sbClient.close();
4. 运行流程
- 启动目标Node API服务(gitWrite所在服务)
- 启动修改后的Service Bus接收端服务,保持后台运行
- 启动后端中转服务(Express服务)
- 启动React前端应用,提交表单触发整个流程
内容的提问来源于stack exchange,提问作者Karan S
相关产品推荐
相关产品推荐

