Python操作Cosmos DB Upsert Item遇BadRequest,如何修正?
问题描述
我编写了Python Azure Functions代码,通过Blob触发器读取文件并操作Azure Cosmos DB,代码如下:
import logging import json import os from azure.cosmos import CosmosClient import azure.functions as func url = os.environ["ACCOUNT_URI"] key = os.environ["ACCOUNT_KEY"] client1 = CosmosClient(url, key) client = CosmosClient.from_connection_string(os.environ["CosmosDBConnStr"]) database_name = "" container_name = "" database = client.get_database_client(database_name) container = database.get_container_client(container_name) logging.info(f"container: {url}") 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") #reading file from blob contract_data=myblob.read() try: logging.info(f"contract data: {contract_data}") contract_json = json.loads(contract_data) version = contract_json.get("version") name = contract_json.get("name") title = contract_json.get("title") logging.info(f"contract json: {contract_json}") query = "SELECT c.version,c.name,c.title,c.Theme,c.description,c['data owner'],c.confidentiality,c.table1 FROM c " items = list(container.query_items( query=query, enable_cross_partition_query=True )) logging.info(f"item: {items[0]}") for item in items: if item["name"] == name and item["version"] == version: if item["title"] == title: logging.info(f"Skipping, item already exists: {item}") return # Skip saving the document container.upsert_item(body=contract_json,pre_trigger_include = None,post_trigger_include= None) return doc.set(func.Document.from_json(contract_data)) except Exception as e: logging.info(f"Error: {e}")
当前问题:当上传同version、同name但title不同的文档时,调用upsert_item出现BadRequest错误,错误信息:Code: BadRequest,Message: {"Errors":["One of the specified inputs is invalid"]}
现有文档示例:
{ "version": "V1", "name": "demo_contract2", "title": "title122", "Theme": "Theme12", "description": "test data contract management2", "data owner": "j.jansen@amsterdam.nl", "confidentiality": "open", "table1": { "description:": "testen", "attribute_1": { "type": "int", "description:": "testen", "identifiability": "identifiable" } } }
待上传文档示例:
{ "version": "V1", "name": "demo_contract2", "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" } } }
需求:调整upsert_item逻辑,实现替换同version和name但title不同的文档。
解决方案
错误原因分析
- 缺少文档ID:Cosmos DB中每个文档必须包含唯一
id字段,upsert_item依赖id+分区键匹配现有文档。若新文档未指定id,系统会尝试生成新ID,但如果容器设置了name+version的唯一键约束,会因唯一键冲突触发BadRequest。 - 全表查询低效且易出错:原代码查询所有文档后再过滤,不仅性能差,还可能因数据量过大导致内存问题。
- 逻辑漏洞:找到匹配文档后直接
return,未处理多个匹配项(虽唯一键约束下不会出现,但逻辑不严谨)。
代码调整要点
- 使用带过滤条件的参数化查询:仅查询同
name和version的文档,提升效率并避免SQL注入。 - 继承现有文档的ID:找到匹配文档后,将其
id赋值给新文档,确保upsert_item能定位到目标文档进行替换。 - 移除冗余的Out绑定调用:直接使用SDK的
upsert_item即可,无需同时使用func.Out[func.Document]绑定。 - 添加分区键检查:确保新文档包含容器的分区键字段(若分区键不是
id)。
调整后的完整代码
import logging import json import os from azure.cosmos import CosmosClient import azure.functions as func # 初始化Cosmos客户端(保留一种初始化方式即可,这里保留连接字符串方式) client = CosmosClient.from_connection_string(os.environ["CosmosDBConnStr"]) database_name = "" container_name = "" database = client.get_database_client(database_name) container = database.get_container_client(container_name) def main(myblob: func.InputStream): logging.info(f"Python blob trigger function processed blob \nName: {myblob.name}") try: # 读取并解析Blob内容 contract_data = myblob.read() contract_json = json.loads(contract_data) version = contract_json.get("version") name = contract_json.get("name") title = contract_json.get("title") if not version or not name: logging.error("Missing required fields: version or name") return # 参数化查询,仅匹配同name和version的文档 query = """ SELECT * FROM c WHERE c.name = @name AND c.version = @version """ items = list(container.query_items( query=query, parameters=[ {"name": "@name", "value": name}, {"name": "@version", "value": version} ], enable_cross_partition_query=True )) if items: existing_item = items[0] if existing_item["title"] == title: logging.info(f"Skipping, item already exists: {existing_item['id']}") return # 继承现有文档的ID和分区键(如果分区键不是id,需额外赋值) contract_json["id"] = existing_item["id"] # 如果容器分区键不是id,例如分区键是"name",则确保contract_json["name"]已存在(这里已有) # 执行upsert替换现有文档 updated_item = container.upsert_item(body=contract_json) logging.info(f"Updated existing item: {updated_item['id']}") else: # 无匹配文档,插入新文档(系统会自动生成id,或可手动指定) new_item = container.create_item(body=contract_json) logging.info(f"Created new item: {new_item['id']}") except Exception as e: logging.error(f"Error processing blob: {str(e)}")
额外检查项
- 容器唯一键配置:若需确保
name+version唯一,需在Cosmos DB容器中设置唯一键策略为/name,/version,避免重复数据。 - 分区键设置:确认容器的分区键字段,确保新文档包含该字段的有效值,否则会触发分区键错误。
- 权限配置:确保Azure Functions的身份拥有Cosmos DB的
Write权限,避免权限不足导致的错误。
内容的提问来源于stack exchange,提问作者aca1803
相关产品推荐
相关产品推荐

