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

如何使用NodeJS读取Azure Event Grid主题中的数据

使用Node.js读取Azure Event Grid主题数据的方法

Azure Event Grid是推送式架构,无法直接从主题中拉取数据,需要先创建主题订阅,将事件推送到指定接收端点,再用Node.js处理这些推送过来的事件。以下是具体实现方案:

1. 创建Event Grid主题订阅

首先需要在Azure门户或通过Azure CLI创建主题订阅,指定事件的接收端点,常用类型包括:

  • Webhook(最常用,可通过Node.js搭建HTTP服务作为端点)
  • Azure Storage Queue
  • Azure Function

这里以Webhook端点为例展开说明。

2. 用Node.js搭建Webhook接收服务

步骤1:初始化项目并安装依赖

创建项目文件夹后,执行以下命令:

npm init -y
npm install express body-parser

步骤2:编写事件接收代码

创建server.js文件,内容如下:

const express = require('express');
const bodyParser = require('body-parser');

const app = express();
app.use(bodyParser.json());

// 处理Event Grid的订阅验证请求(必须实现,否则订阅无法激活)
app.post('/eventgrid-endpoint', (req, res) => {
    const isValidationReq = req.headers['aeg-event-type'] === 'SubscriptionValidation';
    const validationCode = isValidationReq ? req.body[0].data.validationCode : null;

    if (validationCode) {
        res.status(200).json({ validationResponse: validationCode });
        console.log('订阅验证通过');
        return;
    }

    // 处理实际推送的业务事件
    const events = req.body;
    events.forEach(event => {
        console.log('收到事件详情:');
        console.log(`事件ID: ${event.id}`);
        console.log(`事件类型: ${event.eventType}`);
        console.log(`来源主题: ${event.topic}`);
        console.log(`事件数据: ${JSON.stringify(event.data, null, 2)}`);
        console.log('---');
    });

    res.status(200).send('事件接收成功');
});

const PORT = process.env.PORT || 3000;
app.listen(PORT, () => {
    console.log(`服务运行在端口 ${PORT},接收端点:/eventgrid-endpoint`);
});

步骤3:启动服务

执行命令启动HTTP服务:

node server.js

3. 配置主题订阅指向Webhook端点

在Azure门户找到你的Event Grid主题,创建新订阅:

  • 选择“Webhook”作为端点类型
  • 输入服务的公网访问地址(本地测试可使用ngrok等工具暴露端口,格式如https://your-ngrok-id.ngrok.io/eventgrid-endpoint)
  • 完成订阅创建后,服务会自动响应验证请求,订阅激活后,主题中的新事件将实时推送到你的服务。

4. 基于Azure Storage Queue的接收方案

如果选择Storage Queue作为端点,可使用@azure/storage-queue包拉取队列中的事件:

const { QueueClient } = require('@azure/storage-queue');

const storageConnStr = '你的存储账户连接字符串';
const queueName = '你的队列名称';

async function fetchEventsFromQueue() {
    const queueClient = new QueueClient(storageConnStr, queueName);
    const msgResponse = await queueClient.receiveMessages({ numberOfMessages: 10 });
    
    for (const msg of msgResponse.receivedMessages) {
        const eventData = JSON.parse(msg.messageText);
        console.log('从队列获取事件:', eventData);
        // 处理完成后删除队列消息
        await queueClient.deleteMessage(msg.messageId, msg.popReceipt);
    }
}

fetchEventsFromQueue().catch(console.error);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 16:55:23