【问题标题】:APScheduler add lots of jobs concurrently (database Jobstore)APScheduler 同时添加大量作业(数据库 Jobstore)
【发布时间】:2020-02-05 17:47:11
【问题描述】:

如何同时安排大量 APScheduler 作业 (4,000+)? (我必须在某些用户事件之后安排所有这些。)

迭代地调用add_job 对许多工作来说耗时太长。但是当我尝试使用AsyncIOScheduler 和以下异步代码时,我也没有得到任何额外的性能提升。

注意:我的调度程序需要通过 SqlAlchemy 连接到 SQL 作业存储

scheduler = AsyncIOScheduler(jobstores={"default": SQLAlchemyJobStore(url="a valid db connection str")})
scheduler.start()

def schedule_jobs_quickly():
    # init lots of (fake) jobs
    jobs = []
    for i in range(3000):
        jobs.append(i)
    send_time = datetime.datetime.now() + datetime.timedelta(days=2)

    # try to schedule jobs concurrently
    start_time = time.time()
    asyncio.get_event_loop().run_until_complete(schedule_all_jobs(jobs, send_time))
    duration = time.time() - start_time
    print(f"Created {len(jobs)} jobs in {duration} seconds")


async def schedule_all_jobs(all_jobs, send_time):
    tasks = []
    for job in all_jobs:
        task = asyncio.ensure_future(schedule_job(job, send_time))
        tasks.append(task)
    await asyncio.gather(*tasks, return_exceptions=True)


async def schedule_job(job, send_time):
    scheduler.add_job(send_email_if_needed, trigger=send_time)

结果很慢。如何加快速度?

>>> schedule_jobs_quickly()
...
Created 3000 jobs in 401.9982771873474 seconds

作为比较,这是使用默认内存作业存储的BackgroundScheduler() 所花费的时间:

Created 3000 jobs in 0.9155495166778564 seconds

所以,似乎是数据库连接如此昂贵。也许有一种方法可以使用同一连接创建多个作业,而不是为每个 add_job 重新连接?

【问题讨论】:

    标签: python-3.x concurrency python-asyncio apscheduler


    【解决方案1】:

    这不是我正在寻找的解决方案,但我决定放弃 AsyncIOScheduler,而是将我的许多任务安排在一个单独的线程中,这样我的程序的其余部分就可以继续进行,而不会被所有数据库连接所阻碍。下面的例子。

    from threading import Thread
    
    def schedule_jobs_quickly():
        # init lots of (fake) jobs
        jobs = []
        for i in range(3000):
            jobs.append(i)
        send_time = datetime.datetime.now() + datetime.timedelta(days=2)
    
        # schedule jobs in new thread
        scheduler_thread = Thread(target=schedule_email_jobs, args=(email_jobs,))
        scheduler_thread.start()
    
    
    def schedule_email_jobs(jobs):
        for job in jobs:
            scheduler.add_job(send_email, trigger=send_time)
    
    def send_email():
       # sends email 
    

    【讨论】:

    • 你得到了什么样的速度?
    • 不记得了,抱歉
    猜你喜欢
    • 2020-05-20
    • 2017-05-16
    • 1970-01-01
    • 2015-05-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多