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

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不同的文档。


解决方案

错误原因分析

  1. 缺少文档ID:Cosmos DB中每个文档必须包含唯一id字段,upsert_item依赖id+分区键匹配现有文档。若新文档未指定id,系统会尝试生成新ID,但如果容器设置了name+version的唯一键约束,会因唯一键冲突触发BadRequest。
  2. 全表查询低效且易出错:原代码查询所有文档后再过滤,不仅性能差,还可能因数据量过大导致内存问题。
  3. 逻辑漏洞:找到匹配文档后直接return,未处理多个匹配项(虽唯一键约束下不会出现,但逻辑不严谨)。

代码调整要点

  1. 使用带过滤条件的参数化查询:仅查询同name和version的文档,提升效率并避免SQL注入。
  2. 继承现有文档的ID:找到匹配文档后,将其id赋值给新文档,确保upsert_item能定位到目标文档进行替换。
  3. 移除冗余的Out绑定调用:直接使用SDK的upsert_item即可,无需同时使用func.Out[func.Document]绑定。
  4. 添加分区键检查:确保新文档包含容器的分区键字段(若分区键不是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)}")

额外检查项

  1. 容器唯一键配置:若需确保name+version唯一,需在Cosmos DB容器中设置唯一键策略为/name,/version,避免重复数据。
  2. 分区键设置:确认容器的分区键字段,确保新文档包含该字段的有效值,否则会触发分区键错误。
  3. 权限配置:确保Azure Functions的身份拥有Cosmos DB的Write权限,避免权限不足导致的错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:55:05