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

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路径、脚本路径,最后是日志输出(方便排查问题)。


注意事项

  1. 错误处理:上面的代码加了基础的try-except,你可以根据实际需求扩展,比如API请求重试、MongoDB重连逻辑。
  2. 索引验证:可以登录MongoDB客户端,用db.campaignzonestats.getIndexes()(集合名是模型类名小写)验证复合索引是否生效。
  3. 数据一致性:如果API返回的是增量数据,你可能需要用$inc来累加数值(比如"$inc": {"impressions": api_data_item["impressions"]}),而不是直接覆盖——这取决于你的API数据是全量还是增量,根据实际情况调整update逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:50:46