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

如何在消息发送至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. 运行流程

  1. 启动目标Node API服务(gitWrite所在服务)
  2. 启动修改后的Service Bus接收端服务,保持后台运行
  3. 启动后端中转服务(Express服务)
  4. 启动React前端应用,提交表单触发整个流程

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 18:15:27