MongoEngine中Upsert与嵌入文档实现:定时API数据同步需求
嘿,这需求不难搞定!我来帮你一步步实现按campaign_id+zone_id组合UPSERT和每小时自动运行的功能,结合你用的Python3+MongoEngine来搞:
第一步:实现MongoEngine的UPSERT逻辑
首先,你需要确保你的MongoEngine文档模型里,campaign_id和zone_id是复合唯一索引——这既保证了组合的唯一性,也能让UPSERT的查询效率更高。假设你的数据结构包含impressions、clicks这些需要更新的字段,模型可以这么定义:
from mongoengine import Document, StringField, IntField, DateTimeField, Index from datetime import datetime class CampaignZoneStats(Document): campaign_id = StringField(required=True) zone_id = StringField(required=True) impressions = IntField(default=0) clicks = IntField(default=0) # 其他字段根据你的实际数据补充,比如cost、conversions等 updated_at = DateTimeField(default=datetime.utcnow) # 核心:建立campaign_id+zone_id的复合唯一索引 meta = { 'indexes': [ Index('campaign_id', 'zone_id', unique=True) ] }
接下来,把原来的“插入”逻辑改成UPSERT——用MongoEngine的update_one()方法,配合upsert=True参数,就能实现「找到匹配组合就更新,找不到就插入」的效果:
def upsert_single_stat(api_data_item): # 构造查询条件:匹配campaign_id和zone_id的组合 query_filter = { "campaign_id": api_data_item["campaign_id"], "zone_id": api_data_item["zone_id"] } # 构造要更新的字段:这里把API返回的最新数据覆盖/更新到数据库 update_fields = { "$set": { "impressions": api_data_item["impressions"], "clicks": api_data_item["clicks"], # 其他需要同步的字段都放在这里 "updated_at": datetime.utcnow() } } # 执行UPSERT CampaignZoneStats.objects(**query_filter).update_one(**update_fields, upsert=True)
然后在你的数据拉取逻辑里,遍历API返回的每条数据,调用这个函数就行:
def fetch_data_from_api(): # 这里是你原来的API拉取逻辑,返回列表格式的原始数据 # 示例:return [{"campaign_id": "xxx", "zone_id": "yyy", ...}, ...] pass def batch_upsert(): try: api_data = fetch_data_from_api() for item in api_data: upsert_single_stat(item) print(f"Batch upsert finished at {datetime.utcnow()}") except Exception as e: print(f"Error during batch upsert: {str(e)}")
第二步:实现每小时自动运行
有两种常用方案,选适合你的就行:
方案1:用APScheduler(Python代码内集成定时任务)
这个方案适合需要脚本持续运行的场景,比如你的服务器一直开着,脚本后台运行。
首先安装APScheduler:
pip install apscheduler
然后在脚本末尾添加定时任务逻辑:
from apscheduler.schedulers.blocking import BlockingScheduler if __name__ == "__main__": # 初始化调度器,可指定时区(比如国内用Asia/Shanghai) scheduler = BlockingScheduler(timezone="Asia/Shanghai") # 添加每小时运行一次的任务 scheduler.add_job(batch_upsert, "interval", hours=1) print("Scheduler started! Press Ctrl+C to stop.") try: scheduler.start() except KeyboardInterrupt: scheduler.shutdown() print("Scheduler stopped.")
方案2:用系统Crontab(适合Linux服务器)
如果你的脚本不需要一直运行,用系统的定时任务更轻量。
编辑crontab:
crontab -e
添加一行,指定每小时整点运行你的脚本:
0 * * * * /usr/bin/python3 /path/to/your/script.py >> /path/to/logfile.log 2>&1
解释:0 * * * * 表示每小时第0分钟运行,后面是Python路径、脚本路径,最后是日志输出(方便排查问题)。
注意事项
- 错误处理:上面的代码加了基础的try-except,你可以根据实际需求扩展,比如API请求重试、MongoDB重连逻辑。
- 索引验证:可以登录MongoDB客户端,用
db.campaignzonestats.getIndexes()(集合名是模型类名小写)验证复合索引是否生效。 - 数据一致性:如果API返回的是增量数据,你可能需要用
$inc来累加数值(比如"$inc": {"impressions": api_data_item["impressions"]}),而不是直接覆盖——这取决于你的API数据是全量还是增量,根据实际情况调整update逻辑。
内容的提问来源于stack exchange,提问作者Even Even
相关产品推荐
相关产品推荐

