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

Azure PubSub监听Serverless函数本地正常,部署至Azure后写入Mongo失败

问题描述

我有一个网站应用,会生成包含4个属性的简单对象,提交至Azure无服务器函数A。函数A接收请求后,将对象发送至Azure上的Azure PubSub实例。另外还有两个无服务器函数:

  • 函数B:已订阅该PubSub,收到对象后写入Cosmos MongoDB
  • 函数C:供网站应用从MongoDB读取数据

网站、函数A、函数C及PubSub均部署在Azure中。本地运行函数B时整个流程完全正常,但部署到Azure后,函数B能接收PubSub消息,却无法写入MongoDB。已将Azure中函数B的CORS设置为*,怀疑要么是函数B部署时缺少必要配置,要么是本地允许写入但Azure环境不允许,想确认Cosmos是否有类似MongoDB Atlas的客户端IP允许设置。


相关代码文件

函数B代码

module.exports = function(context, req) {
    const WebSocket = require('ws');
    const { WebPubSubServiceClient } = require('@azure/web-pubsub');
    const MongoClient = require('mongodb').MongoClient;
    
  
    const connect = async () => {
      client = await MongoClient.connect('connection string' );
                                          
      return client;
      }
  
    main();
  
    async function main() {
      const hub = "kurtpubsub.webpubsub.azure.com";
      var connectionString = "Endpoint=https://kurtpubsub.webpubsub.azure.com;AccessKey= key";
      let serviceClient = new WebPubSubServiceClient(connectionString, hub);
      let token = await serviceClient.getClientAccessToken();
      let ws = new WebSocket(token.url);
      ws.on('open', () => console.log('connected'));
      ws.on('message', data => {

       let newRecord = JSON.parse(data);
        // got the data, now write to mongo
      
        // have to wait for 2nd promise to get he client object
        connect().then( console.log("connected to mongo"))
        .then((client) => { writeDB(client, newRecord)});
        //console.log(newRecord);
      });
    }
  
    function writeDB(client, pNewRecord)
    {
      try {
        const database = client.db("restdb");
        //context.log("connected to DB successfully");
          // Insert a single document
          
          database.collection('restaurants').insertOne(pNewRecord);
          //client.close();
        } catch (err) {
          console.log(err.stack);
          client.close();     // Close connection if app is dying
        }
    }
  }

host.json文件

{
  "version": "2.0",
  "logging": {
    "applicationInsights": {
      "samplingSettings": {
        "isEnabled": true,
        "excludedTypes": "Request"
      }
    }
  },
  "extensionBundle": {
    "id": "Microsoft.Azure.Functions.ExtensionBundle",
    "version": "[3.*, 4.0.0)"
  }
}

function.json文件

{
  "bindings": [
    {
      "authLevel": "anonymous",
      "type": "httpTrigger",
      "direction": "in",
      "name": "req",
      "methods": [
        "get",
        "post"
      ]
    },
    {
      "type": "http",
      "direction": "out",
      "name": "res"
    }
  ]
}

问题分析与解决方案

1. Cosmos DB的IP访问控制配置

Cosmos DB确实有类似MongoDB Atlas的IP白名单机制,这是最可能的诱因:

  • 登录Azure门户,找到目标Cosmos MongoDB账户
  • 进入网络>防火墙和虚拟网络设置:
    • 如果设置为「允许从选定的网络访问」,需将函数B的出站IP地址添加到白名单(可在函数应用>设置>属性中查看出站IP)
    • 临时勾选「允许从所有网络访问」测试是否解决问题(测试后务必改回安全配置)
    • 确保勾选「允许Azure服务和资源访问此账户」,让Azure内部服务(如函数B)能正常访问Cosmos DB

2. 函数B代码的异步逻辑问题

代码存在异步处理缺陷,本地运行可能因时序宽松掩盖问题,部署到Azure后会触发写入失败:

  • connect().then( console.log("connected to mongo"))中console.log会立即执行,并非连接成功后触发,应改为connect().then(() => console.log("connected to mongo"))
  • insertOne是异步操作,未等待执行完成也未处理错误,Azure函数可能因进程提前结束导致写入中断
  • 每次收到消息都创建新Mongo连接,易引发连接池耗尽,应复用连接

修改后的优化代码:

module.exports = async function(context, req) {
    const { WebPubSubServiceClient } = require('@azure/web-pubsub');
    const { MongoClient } = require('mongodb');

    // 复用Mongo连接,避免重复创建
    let mongoClient;
    async function connectMongo() {
        if (!mongoClient || !mongoClient.isConnected()) {
            mongoClient = await MongoClient.connect(process.env.MONGO_CONNECTION_STRING);
            context.log('MongoDB连接成功');
        }
        return mongoClient;
    }

    async function main() {
        const hub = "kurtpubsub"; // hub应为PubSub实例中的hub名称,而非完整域名
        const connectionString = process.env.PUBSUB_CONNECTION_STRING;
        const serviceClient = new WebPubSubServiceClient(connectionString, hub);
        const token = await serviceClient.getClientAccessToken();
        
        const WebSocket = require('ws');
        const ws = new WebSocket(token.url);
        
        ws.on('open', () => context.log('已连接到PubSub'));
        ws.on('message', async (data) => {
            try {
                const newRecord = JSON.parse(data);
                context.log('收到PubSub消息:', newRecord);
                
                const client = await connectMongo();
                const database = client.db("restdb");
                // 等待写入完成并记录结果
                const result = await database.collection('restaurants').insertOne(newRecord);
                context.log('写入MongoDB成功,ID:', result.insertedId);
            } catch (err) {
                context.error('处理消息失败:', err.stack);
                // 异常时关闭连接并重置
                if (mongoClient) {
                    await mongoClient.close();
                    mongoClient = null;
                }
            }
        });
        ws.on('error', (err) => context.error('WebSocket错误:', err));
    }

    await main();
    context.res = { status: 200, body: '函数已启动并监听PubSub' };
}

3. Azure函数配置检查

  • 不要硬编码连接字符串,将MongoDB和PubSub的连接字符串存入函数应用的应用设置中,通过process.env.变量名读取
  • 查看函数B的运行日志(函数应用>监控>日志),获取MongoDB连接或写入的具体错误信息,这是排查核心
  • 当前function.json使用HTTP触发器,但函数B作为PubSub监听器,更推荐使用Web PubSub触发器绑定,无需手动管理WebSocket连接,稳定性更高

4. Web PubSub参数修正

代码中hub参数应为PubSub实例内创建的hub名称,而非完整域名(例如kurtpubsub而非kurtpubsub.webpubsub.azure.com)


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 05:13:21