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

当Blob文件实体变更时,如何将其作为新文档存入Azure Cosmos DB?

问题描述

现有Python Azure函数代码仅在Blob文件的name和version字段与Cosmos DB中现有文档都不同时,才将Blob内容存入Cosmos DB的dmdb数据库DataContract容器。现在需要调整逻辑:当Blob文件内除name外的其他字段(如data owner、table1的描述等)发生变更时,即使version相同,也将该文件作为新文档存入Cosmos DB。

现有函数代码

import logging
import json
import os
from azure.cosmos import CosmosClient
import azure.functions as func

client = CosmosClient.from_connection_string(os.environ["CosmosDBConnStr"])

database_name = "dmdb"
container_name = "DataContract"

database = client.get_database_client(database_name)
container = database.get_container_client(container_name)
logging.info(f"container: {container}")
def main(myblob: func.InputStream, doc: func.Out[func.Document]):
    logging.info(f"Python blob trigger function processed blob \n"
                 f"Name: {myblob.name}\n")


contract_data=myblob.read()
blob_filename = os.path.basename(myblob.name)
logging.info(f"file name: {blob_filename}")
logging.info(f"my blob: {myblob}")
logging.info(f"my document: {doc}")
doc_value = doc.get()
logging.info(f"my document1: {doc_value}")

try:
    logging.info(f"contract data: {contract_data}")
    contract_json = json.loads(contract_data)
    version = contract_json.get("version")
    name = contract_json.get("name")
    query = f''' SELECT c.name,c.version
            FROM c 
            '''
    items = list(container.query_items(query, enable_cross_partition_query=True))
    logging.info(f"items: {items}")

    for item in items:
        if item["name"] == name and item["version"] == version:
            logging.info(f"Skipping, item already exists: {item}")
            return
    doc.set(func.Document.from_json(contract_data))
    doc_value = doc.get()
    version = doc_value.get("version")
    name = doc_value.get("name")
    logging.info(f"Version: {version}")
    logging.info(f"Name: {name}")

except Exception as e:
        logging.info(f"Error: {e}")

示例数据

Blob文件中的contract_data示例

{
    "version": "V3",
    "name": "demo_contract9",
    "title": "title1",
    "Theme": "Theme1",
    "description": "test data contract management2",
    "data owner": "aysegul@amsterdam.nl",
    "confidentiality": "open",
    "table1": {
        "description:": "test",
        "attribute_1": {
            "type": "int",
            "description:": "testen",
            "identifiability": "identifiable"
        }
    }
}  

Cosmos DB中已有文档示例

{
    "version": "V2",
    "name": "demo_contract9",
    "title": "title1",
    "Theme": "Theme1",
    "description": "test data contract management2",
    "data owner": "j.jansen@amsterdam.nl",
    "confidentiality": "open",
    "table1": {
        "description:": "testen",
        "attribute_1": {
            "type": "int",
            "description:": "testen",
            "identifiability": "identifiable"
        }
    },
    "id": "d4923d64-96df-4de4-baf8-a3e902c86528"
}
实现方案

要实现需求,核心是对比Blob内容与Cosmos DB中同name的所有文档的完整内容(排除自动生成的id字段),只要存在差异就插入新文档。具体修改如下:

关键修改点

  1. 调整查询语句:查询同name的所有完整文档,而非仅name和version字段
  2. 对比逻辑:排除Cosmos DB自动生成的id字段后,对比Blob的JSON数据与现有文档的内容
  3. 移除原有的version相同即跳过的逻辑,改为内容完全一致才跳过

修改后的完整代码

import logging
import json
import os
from azure.cosmos import CosmosClient
import azure.functions as func

client = CosmosClient.from_connection_string(os.environ["CosmosDBConnStr"])

database_name = "dmdb"
container_name = "DataContract"

database = client.get_database_client(database_name)
container = database.get_container_client(container_name)
logging.info(f"container: {container}")

def main(myblob: func.InputStream, doc: func.Out[func.Document]):
    logging.info(f"Python blob trigger function processed blob \n"
                 f"Name: {myblob.name}\n")

    contract_data = myblob.read()
    blob_filename = os.path.basename(myblob.name)
    logging.info(f"file name: {blob_filename}")
    
    try:
        logging.info(f"contract data: {contract_data}")
        contract_json = json.loads(contract_data)
        name = contract_json.get("name")
        
        # 查询所有同name的完整文档,使用参数化查询避免SQL注入
        query = '''
            SELECT *
            FROM c 
            WHERE c.name = @name
        '''
        parameters = [{"name": "@name", "value": name}]
        items = list(container.query_items(
            query=query,
            parameters=parameters,
            enable_cross_partition_query=True
        ))
        logging.info(f"Found {len(items)} existing items with name: {name}")

        # 标记是否存在完全一致的文档(排除自动生成的id)
        exists_identical = False
        for item in items:
            # 复制现有文档并移除id字段,用于对比
            item_without_id = item.copy()
            del item_without_id["id"]
            # 对比内容是否完全一致
            if item_without_id == contract_json:
                exists_identical = True
                logging.info(f"Skipping, identical document already exists: {item['id']}")
                break
        
        # 如果没有完全一致的文档,插入新文档
        if not exists_identical:
            doc.set(func.Document.from_json(contract_data))
            new_doc = doc.get()
            logging.info(f"Inserted new document. Name: {new_doc.get('name')}, Version: {new_doc.get('version')}")

    except Exception as e:
        logging.error(f"Error: {str(e)}")

代码说明

  • 查询逻辑:使用参数化查询获取所有同name的文档,避免SQL注入风险
  • 内容对比:移除Cosmos DB自动生成的id字段(每个新文档的id唯一,不参与内容对比),与Blob的JSON数据做完全匹配
  • 插入逻辑:仅当没有完全一致的文档时才插入新文档,确保任何内容变更都会生成新的Cosmos DB文档

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 23:27:03