【问题标题】:Increment a counter and trigger an action when a threshold is exceeded增加一个计数器并在超过阈值时触发一个动作
【发布时间】:2023-04-10 01:27:01
【问题描述】:

我有一个这样的模型

class Thingy(models.Model):
    # ...
    failures_count = models.IntegerField()

我有需要执行此操作的并发进程(Celery 任务):

  1. 做一些处理
  2. 如果处理失败,则增加相应failures_counterThingy
  3. 如果failures_counter 超过某些Thingy 的阈值,则发出警告,但只有一个警告。

我对如何在没有竞争条件的情况下执行此操作有一些想法,例如使用显式锁(通过select_for_update):

@transaction.commit_on_success
def report_failure(thingy_id):
    current, = (Thingy.objects
               .select_for_update()
               .filter(id=thingy_id)
               .values_list('failures_count'))[0]
    if current == THRESHOLD:
        issue_warning_for(thingy_id)
    Thingy.objects.filter(id=thingy_id).update(
        failures_count=F('failures_count') + 1
    )

或者通过使用 Redis(它已经存在)进行同步:

@transaction.commit_on_success
def report_failure(thingy_id):
    Thingy.objects.filter(id=thingy_id).update(
        failures_count=F('failures_count') + 1
    )
    value = Thingy.objects.get(id=thingy_id).only('failures_count').failures_count
    if value >= THRESHOLD:
        if redis.incr('issued_warning_%s' % thingy_id) == 1:
            issue_warning_for(thingy_id)

两种解决方案都使用锁。由于我使用的是 PostgreSQL,有没有办法在不锁定的情况下实现这一点?


我正在编辑问题以包含 答案(感谢 Sean Vieira,请参阅下面的答案)。该问题询问了一种避免锁定的方法,这个答案是最佳的,因为它利用了multi-version concurrency control (MVCC) as implemented by PostgreSQL

这个特定问题明确允许使用 PostgreSQL 功能,尽管许多 RDBMS 实现了UPDATE ... RETURNING,但它不是标准 SQL,并且 Django 的 ORM 不支持开箱即用,因此它需要通过 raw() 使用原始 SQL。相同的 SQL 语句将在其他 RDBMS 中工作,但每个引擎都需要自己讨论同步、事务隔离和并发模型(例如,带有 MyISAM 的 MySQL 仍将使用锁)。

def report_failure(thingy_id):
    with transaction.commit_on_success():
        failure_count = Thingy.objects.raw("""
            UPDATE Thingy
            SET failure_count = failure_count + 1
            WHERE id = %s
            RETURNING failure_count;
        """, [thingy_id])[0].failure_count

    if failure_count == THRESHOLD:
        issue_warning_for(thingy_id)

【问题讨论】:

  • 最简单的方法是在 redis 中拥有两个计数器...
  • @armonge 这对于与我正在使用的设置不同的设置很有用。在我的设置中,我需要长期存储故障计数,而 Redis 仅用于缓存/同步。

标签: python django postgresql concurrency


【解决方案1】:

据我所知,Django 的 ORM 不支持这个开箱即用 - 然而,这并不意味着它不能完成,你只需要深入到 SQL 级别(暴露,在 Django 的ORM 通过Managerraw method) 使其工作。

如果您使用 PostgresSQL >= 8.2,那么您可以使用 RETURNING 来获取 failure_count 的最终值,而无需任何额外的锁定(数据库仍将锁定,但只有足够长的时间来设置值,无需额外的时间失去与你的沟通):

# ASSUMPTIONS: All IDs are valid and IDs are unique
# More defenses are necessary if either of these assumptions
# are not true.
failure_count = Thingy.objects.raw("""
    UPDATE Thingy
    SET failure_count = failure_count + 1
    WHERE id = %s
    RETURNING failure_count;
""", [thingy_id])[0].failure_count

if failure_count == THRESHOLD:
    issue_warning_for(thingy_id)

【讨论】:

  • Sean,我知道UPDATE RETURNING,它看起来很整洁,但我不确定它是否解决了 Django 默认事务隔离级别 (READ COMMITTED) 的问题。如果两个并发事务发出相同的语句会发生什么?他们会得到两个不同的值吗?如果其中一项事务回滚怎么办?
  • @DavideR。 - 这正是 PostGresSQL 为您处理的事情 - concurrent updates to the same data will happen serially。如果一项交易失败,世界仍将保持一致状态。
  • 谢谢,这个链接很好地阐明了它。对于所有感兴趣的人,幻灯片 11 和 14 说明了 READ COMMITTED(Django 的默认设置)的行为,在这种情况下。还需要注意的是,PostgreSQL MVCC 在任何情况下都保证了这种行为,并且在不使用锁的情况下实现了高并发。
【解决方案2】:

我真的不知道你必须在没有锁定的情况下完成这项工作的原因,你有多少个任务同时运行?

但是,我认为有一种方法可以做到这一点,而无需像这样锁定:

你应该有另一个模型,例如失败:

class Failure(models.Model):
    thingy = models.ForeignKey(Thingy)

您的 *report_failure* 应该是这样的:

from django.db import transaction
@transaction.commit_manually
def flush_transaction():
    transaction.commit()

@transaction.commit_on_success
def report_failure(thingy_id):
    thingy = Thingy.objects.get(id=thingy_id)
    #uncomment following line if you found that the query is cached (not get updated result)
    #flush_transaction()

    current = thingy.failure_set.count()
    if current >= THRESHOLD:
        issue_warning_for(thingy_id)
    Failure.objects.create(thingy=thingy)

我知道这种方法很糟糕,因为它会创建很多失败记录。但这是我能想出的唯一想法。对此感到抱歉。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-16
    • 2018-12-31
    • 1970-01-01
    • 2019-11-15
    • 1970-01-01
    • 2015-08-08
    相关资源
    最近更新 更多