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

Timer Trigger函数集成CosmosDB异常:无法执行数据库操作

Python Timer Trigger函数集成CosmosDB无执行无报错问题排查解决

问题描述

我有一个Timer Trigger函数应用,用于与Telegram Bot交互,API请求后发送消息。本地运行函数应用正常,但集成CosmosDB存储数据时出现问题,无法保存信息。已配置好Telegram和CosmosDB的相关变量,使用func host start --port 7072启动函数后,发现print('connected to db')未输出,所有CosmosDB相关操作都未执行,终端无报错信息。

代码片段

try:
    database_obj  = client.get_database_client(database_name)
    await database_obj.read()
    return database_obj
except exceptions.CosmosResourceNotFoundError:
    print("Creating database")
    return await client.create_database(database_name)
# </create_database_if_not_exists>
    
# Create a container
# Using a good partition key improves the performance of database operations.
# <create_container_if_not_exists>
async def get_or_create_container(database_obj, container_name):
    try:        
        todo_items_container = database_obj.get_container_client(container_name)
        await todo_items_container.read()   
        return todo_items_container
    except exceptions.CosmosResourceNotFoundError:
        print("Creating container with lastName as partition key")
        return await database_obj.create_container(
            id=container_name,
            partition_key=PartitionKey(path="/lastName"),
            offer_throughput=400)
    except exceptions.CosmosHttpResponseError:
        raise
# </create_container_if_not_exists>

async def populate_container_items(container_obj, items_to_create):
    # Add items to the container
    family_items_to_create = items_to_create
    # <create_item>
    for family_item in family_items_to_create:
        inserted_item = await container_obj.create_item(body=family_item)
        print("Inserted item for %s family. Item Id: %s" %(inserted_item['lastName'], inserted_item['id']))
    # </create_item>
# </method_populate_container_items>

async def read_items(container_obj, items_to_read):
    # Read items (key value lookups by partition key and id, aka point reads)
    # <read_item>
    for family in items_to_read:
        item_response = await container_obj.read_item(item=family['id'], partition_key=family['lastName'])
        request_charge = container_obj.client_connection.last_response_headers['x-ms-request-charge']
        print('Read item with id {0}. Operation consumed {1} request units'.format(item_response['id'], (request_charge)))
    # </read_item>
# </method_read_items>

# <method_query_items>
async def query_items(container_obj, query_text):
    # enable_cross_partition_query should be set to True as the container is partitioned
    # In this case, we do have to await the asynchronous iterator object since logic
    # within the query_items() method makes network calls to verify the partition key
    # definition in the container
    # <query_items>
    query_items_response = container_obj.query_items(
        query=query_text,
        enable_cross_partition_query=True
    )
    request_charge = container_obj.client_connection.last_response_headers['x-ms-request-charge']
    items = [item async for item in query_items_response]
    print('Query returned {0} items. Operation consumed {1} request units'.format(len(items), request_charge))
    # </query_items>
# </method_query_items>

async def run_sample():
    print('aaaa')
    print('sss {0}'.format(CosmosClient(endpoint,credential=key)))
    async with CosmosClient(endpoint, credential = key) as client:
        print('connected to db')
        try:
            database_obj = await get_or_create_db(client, database_name)
            # create a container
            container_obj = await get_or_create_container(database_obj, container_name)
            family_items_to_create = ["link", "ss", "s", "s"]
            await populate_container_items(container_obj, family_items_to_create)
            await read_items(container_obj, family_items_to_create)
            # Query these items using the SQL query syntax. 
            # Specifying the partition key value in the query allows Cosmos DB to retrieve data only from the relevant partitions, which improves performance
            query = "SELECT * FROM c "
            await query_items(container_obj, query)   
        except exceptions.CosmosHttpResponseError as e:
            print('\nrun_sample has caught an error. {0}'.format(e.message))
        finally:
            print("\nQuickstart complete")

async def main(mytimer: func.TimerRequest) -> None:
    utc_timestamp = datetime.datetime.utcnow().replace(
        tzinfo=datetime.timezone.utc).isoformat()
    
    asyncio.create_task(run_sample())
    logging.info(' sono partito')
    sendNews()
    if mytimer.past_due:
        logging.info('The timer is past due!')

    logging.info('Python timer trigger function ran at %s', utc_timestamp)

解决方案

1. 修复异步任务未等待的核心问题

main函数中使用asyncio.create_task(run_sample())仅创建了异步任务,但未等待其完成。Timer Trigger函数执行完同步代码后会直接退出,导致CosmosDB操作还未执行就被终止。将代码修改为:

# 替换asyncio.create_task(run_sample())为:
await run_sample()

2. 修正CosmosDB Item的数据格式

family_items_to_create传入的是字符串列表,但CosmosDB要求Item为包含分区键(lastName)的字典格式。修改示例:

family_items_to_create = [
    {"id": "item_1", "lastName": "link", "message": "sample content 1"},
    {"id": "item_2", "lastName": "ss", "message": "sample content 2"}
]

3. 统一使用日志输出替代print

函数应用环境中print输出可能无法正常显示在终端,改用logging模块确保执行状态可追踪:

# 替换print('connected to db')为:
logging.info('connected to db')

4. 扩大异常捕获范围

当前仅捕获CosmosHttpResponseError,无法覆盖数据格式错误等其他异常。修改run_sample的异常捕获逻辑:

try:
    # ... 原有CosmosDB操作代码
except exceptions.CosmosHttpResponseError as e:
    logging.error(f"CosmosDB请求错误: {e.message}")
except Exception as e:
    logging.error(f"CosmosDB操作失败: {str(e)}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 06:47:03