如何使用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
相关产品推荐
相关产品推荐

