当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字段),只要存在差异就插入新文档。具体修改如下:
关键修改点
- 调整查询语句:查询同
name的所有完整文档,而非仅name和version字段 - 对比逻辑:排除Cosmos DB自动生成的
id字段后,对比Blob的JSON数据与现有文档的内容 - 移除原有的
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
相关产品推荐
相关产品推荐

