【问题标题】:Implementing a Timer in Python在 Python 中实现定时器
【发布时间】:2013-06-05 20:47:44
【问题描述】:

一般概述

  • 我有一个中等规模的 django 项目
  • 我在内存中有一堆前缀树(而不是 DB)
  • 这些树的节点表示受到超时影响的实体/对象。即,我需要在不同的时间点使这些节点超时

设计:

  • 本质上,我需要一个 Timer 构造,它允许我触发一个可重置的 1-shot 计时器并关联并给它一个回调,该回调可以对创建计时器的实体执行一些操作,在本例中是树。

在浏览了各种选项后,我找不到任何可以原生使用的东西(比如一些 django 应用程序)。 Python 中的 Timer 对象不适合这种情况,因为它不会扩展/执行。因此,我决定根据以下内容编写自己的计时器:

  1. 包含时间范围的时间增量对象的排序列表
  2. 触发“tick”的机制

实施选择:

  1. 在 Bisect 周围使用了一个包装器,用于排序的 delta 列表: http://code.activestate.com/recipes/577197-sortedcollection/
  2. 使用 celery 提供滴答声 - 1 分钟的粒度,工作人员将触发我的 Timer 类提供的 timer_tick 函数。 timer_tick 本质上应该通过排序列表,每次滴答都会减少头节点。然后,如果有任何节点的记号为 0,则启动回调并将这些节点从排序的计时器列表中删除。
  3. 创建计时器涉及实例化一个返回对象 ID 的 Timer 对象。此 id 存储在 db 中,并与 DB 中的条目相关联,该条目表示创建计时器的实体

其他数据结构: 为了跟踪 Timer 实例(为每个计时器创建实例化),我有一个 WeakRef 字典,它将 id 映射到 obj

所以本质上,我的主要 Django 项目的内存中有 2 个数据结构。

问题陈述:

由于 celery worker 需要遍历计时器列表并且还可能修改 id2obj 映射,看来我需要找到一种方法在我的 celery worker 和 main 之间共享状态

通过 SO/Google,我发现以下建议

  1. 经理
  2. 共享内存

不幸的是,bisect wrapper 不适合酸洗和/或状态共享。我通过创建一个字典并尝试将排序列表嵌入到字典中来尝试管理器方法。它出现了一个错误(我猜是因为排序列表持有的内存没有共享并将其嵌入到“共享”内存对象将不起作用)

最后...问题:

  1. 有什么方法可以与工作线程共享我的 SortedCollection 和 Weakref Dict

替代解决方案:

如何保持工作线程简单...让它在每个滴答时写入 DB,然后使用 post Db 信号在主线程上获得通知并在主线程中执行过期计时器的处理。当然,缺点是我失去了并行性。

【问题讨论】:

  • 老兄,这是一篇博文还是一个问题?难以置信的大!
  • 旁注:我认为使用搜索树(例如,来自bintreesblist 的内容)而不是带有bisect 的列表会好很多。使用这种设计,每次添加/删除/过期计时器时,都是 O(N) 操作。
  • 更重要的是:“bisect wrapper 不适合酸洗和/或状态共享”是什么意思?它只是一个列表,可能带有一些额外的状态变量,这对于腌制或共享来说是微不足道的——或者更简单的是,使用 multiprocessing.Array 而不是 list
  • 此外,您在谈论线程和进程之间来回切换。他们不是一回事。使用两个线程,默认情况下共享所有内容;您所要做的就是添加适当的同步。对于两个进程,默认情况下不会共享任何内容;您必须明确传递或分享它。
  • 最后,如果你试图在正确隔离的进程之间共享某些东西,并且对于插入、删除、任意查找和查找最低值大约为 O(log N) 或更好……听起来像数据库的工作。有理由不在这里使用吗?

标签: python django timer django-celery celerybeat


【解决方案1】:

让我们从现有实现的一些 cmets 开始:

使用 Bisect 的包装器用于排序的 delta 列表:http://code.activestate.com/recipes/577197-sortedcollection/

虽然这会为您提供 O(1) 次弹出(只要您以相反的时间顺序保留列表),但它会使每个插入 O(N) (同样对于不太常见的操作,例如删除任意作业,如果您有一个 "取消”API)。由于您执行的插入次数与弹出次数完全相同,这意味着整个事情在算法上并不比未排序的列表好。

heapq 替换它(这正是它们的用途)给你 O(log N) 插入。 (注意Python的heapq没有peek,但那是因为heap[0]等价于heap.peek(0),所以你不需要它。)

如果您还需要进行其他操作(取消、非破坏性迭代等)O(log N),您需要一个搜索树;看看 PyPI 上的 blistbintrees 有什么好的。


与 celery 一起提供滴答声 - 1 分钟的粒度,工作人员将触发我的 Timer 类提供的 timer_tick 函数。 timer_tick 本质上应该通过排序列表,每次滴答都会减少头节点。然后,如果有任何节点已被标记为 0,则启动回调并将这些节点从排序的计时器列表中删除。

只保留目标时间而不是增量要好得多。有了目标时间,您只需这样做:

while q.peek().timestamp <= now():
    process(q.pop())

再一次,这是 O(1) 而不是 O(N),而且它要简单得多,并且它将队列中的元素视为不可变的,并且它避免了迭代花费比您的滴答时间更长的任何可能的问题(可能1 分钟的滴答声不是问题……)。


现在,谈谈你的主要问题:

有什么方法可以分享我的 SortedCollection

是的。如果您只想要(timestamp, id) 对的优先级堆,则可以像list 一样轻松地将其放入multiprocessing.Array,除了需要明确跟踪长度。然后你只需要同步每个操作,然后……就是这样。

如果您每分钟只计时一次,并且您预计会更频繁地忙碌,您可以使用 Lock 进行同步,并让 schedule-worker(s) 自动计时。

但老实说,我会完全放弃滴答声,只使用Condition——它更灵活,概念上更简单(即使代码多一点),这意味着当有当您处于负载状态时,无需完成任何工作并快速流畅地响应。例如:

def schedule_job(timestamp, job):
    job_id = add_job_to_shared_dict(job) # see below
    with scheduler_condition:
        scheduler_heap.push((timestamp, job))
        scheduler_condition.notify_all()

def scheduler_worker_run_once():
    with scheduler_condition:
        while True:
            top = scheduler_heap.peek()
            if top is not None:
                delay = top[0] - now()
                if delay <= 0:
                    break
                scheduler_condition.wait(delay)
            else:
                scheduler_condition.wait()
        top = scheduler_heap.pop()
        if top is not None:
            job = pop_job_from_shared_dict(top[1])
            process_job(job)

无论如何,这将我们带到了充满工作的weakdict。

由于weakdict 显式存储对进程内对象的引用,因此跨进程共享它没有任何意义。您要存储的是定义作业实际是什么的不可变对象,而不是可变作业本身。那么它只是一个普通的旧字典。

但是,一个普通的旧字典仍然不是一件容易跨进程共享的事情。

简单的方法是使用dbm 数据库(或shelve 包装器)而不是内存中的dict,与Lock 同步。但这意味着每次有人想要更改数据库时都要重新刷新和重新打开数据库,这可能是不可接受的。

例如,切换到 sqlite3 数据库可能看起来有点矫枉过正,但它可能要简单得多。

另一方面……您实际上在这里进行的唯一操作是“将下一个 id 映射到此作业并返回 id”和“弹出并返回此 id 指定的作业”。这真的需要一个字典吗?键是整数,您可以控制它们。一个Array,加上一个用于下一个键的Value,以及一个Lock,你就差不多完成了。问题是您需要某种密钥溢出方案。不仅是next_id += 1,您还必须滚动并检查已使用的插槽:

with lock:
    next_id += 1
    if next_id == size: next_id = 0
    if arr[next_id] is None:
        arr[next_id] = job
        return next_id

另一种选择是将dict存储在主进程中,并使用Queue让其他进程查询它。

【讨论】:

  • 顺便说一下,虽然定时调度程序往往在事件循环驱动的代码(流服务器、游戏、GUI 应用程序......)中最有用,但你不可能是唯一一个在寻找的人用于多处理(或多线程或greenlet,至少足够相似以适应)代码。你试过在 PyPI 上搜索timerschedule 等吗?
猜你喜欢
  • 2022-08-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多