Python中如何实现Azure Cosmos DB数据存储任务后台执行?
问题根源分析
你的代码无法实现后台任务的核心原因有两个:
storeDataToDb是同步阻塞函数,time.sleep会彻底卡住asyncio事件循环,导致asyncio.run必须等待存储任务完全结束才会返回。- Azure Function的Python运行时在函数返回后会快速回收进程,哪怕你用asyncio启动了后台任务,也会被强制终止,无法完成执行。
可行解决方案
方案一:用Azure Queue Storage解耦任务(最推荐,无状态场景首选)
把数据存储任务拆成两步:原函数只负责生成报告并将任务参数发送到队列,另一个Queue Trigger函数监听队列执行存储逻辑,失败时触发告警。
1. 修改原报告生成函数
import os import json from azure.storage.queue import QueueClient, BinaryBase64EncodePolicy def exectuteReport(): myobj.readInputFile() myobj.readDatafromDb() myobj.generateReport() # 准备要传递给后台任务的数据集 task_payload = json.dumps({"dataset": myobj.dataset}) # 初始化队列客户端(连接字符串从Function应用设置中读取) queue_client = QueueClient.from_connection_string( conn_str=os.getenv("AZURE_STORAGE_CONNECTION_STRING"), queue_name="data-store-task-queue", message_encode_policy=BinaryBase64EncodePolicy() ) # 发送任务到队列 queue_client.send_message(task_payload) print("数据存储任务已提交到后台队列") return {"success": True}
2. 新增Queue Trigger函数(负责存储+告警)
import os import json import time import smtplib from email.mime.text import MIMEText from azure.storage.queue import BinaryBase64DecodePolicy def store_data_trigger(msg): # 解析队列消息 payload = json.loads(msg.get_body(decode_policy=BinaryBase64DecodePolicy())) dataset = payload["dataset"] try: # 执行分批次存储逻辑 batch_count = 0 for batch in dataset: print(f"正在存储第{batch_count}批数据") # 替换为你的Cosmos DB批量写入代码 # cosmos_container.upsert_items(batch) batch_count += 1 time.sleep(2) print("所有数据存储完成") except Exception as e: # 触发告警邮件 send_alert(f"数据存储任务失败,错误信息:{str(e)}") # 抛出异常让Azure Function自动重试(默认最多5次) raise def send_alert(error_content): # 建议用Azure SendGrid服务,以下是SMTP示例 smtp_server = os.getenv("ALERT_SMTP_SERVER") smtp_port = int(os.getenv("ALERT_SMTP_PORT")) sender = os.getenv("ALERT_SENDER_EMAIL") receiver = os.getenv("ALERT_RECEIVER_EMAIL") password = os.getenv("ALERT_EMAIL_PASSWORD") msg = MIMEText(error_content) msg["Subject"] = "Azure Function 数据存储任务告警" msg["From"] = sender msg["To"] = receiver with smtplib.SMTP(smtp_server, smtp_port) as server: server.starttls() server.login(sender, password) server.send_message(msg)
方案二:用Durable Functions(需跟踪任务状态场景)
如果需要监控后台任务的执行状态,可使用Azure Durable Functions,将存储任务封装为Activity函数,原函数作为Orchestrator触发后台任务后立即返回。
1. Orchestrator函数(报告生成+触发后台任务)
import azure.durable_functions as df def orchestrator(context: df.DurableOrchestrationContext): # 执行报告生成逻辑 myobj.readInputFile() myobj.readDatafromDb() myobj.generateReport() # 触发后台存储任务,不等待完成 yield context.call_activity("StoreDataActivity", myobj.dataset) return {"success": True} main = df.Orchestrator.create(orchestrator_function)
2. Activity函数(存储+告警)
import os import time import smtplib from email.mime.text import MIMEText def store_data_activity(dataset): try: batch_count = 0 for batch in dataset: print(f"正在存储第{batch_count}批数据") # Cosmos DB写入逻辑 batch_count += 1 time.sleep(2) print("存储完成") except Exception as e: send_alert(f"Durable任务数据存储失败:{str(e)}") raise def send_alert(error_content): # 同方案一的邮件发送逻辑 smtp_server = os.getenv("ALERT_SMTP_SERVER") smtp_port = int(os.getenv("ALERT_SMTP_PORT")) sender = os.getenv("ALERT_SENDER_EMAIL") receiver = os.getenv("ALERT_RECEIVER_EMAIL") password = os.getenv("ALERT_EMAIL_PASSWORD") msg = MIMEText(error_content) msg["Subject"] = "Durable Function 数据存储告警" msg["From"] = sender msg["To"] = receiver with smtplib.SMTP(smtp_server, smtp_port) as server: server.starttls() server.login(sender, password) server.send_message(msg)
关键注意事项
- 所有敏感信息(存储连接字符串、邮箱密码)必须存在Azure Function的应用设置中,禁止硬编码。
- 使用Queue Storage时,需设置合理的消息可见性超时(建议大于存储任务的最大执行时间),避免消息被重复处理。
- 邮件告警优先使用Azure SendGrid服务,比自建SMTP更稳定且符合Azure最佳实践。
内容的提问来源于stack exchange,提问作者Smith Dwayne
相关产品推荐
相关产品推荐

