☰
Huey 实战配方手册:优雅关闭、多队列、Redis 高可用与 16 个任务队列高级用法
2026/10/10 2:18:07 网站建设 项目流程
  • 任务调度
  • 后端

【免费下载链接】huey

a little task queue for python

项目地址:https://gitcode.com/gh_mirrors/hu/huey
点击查看免费下载

本指南以 huey 官方 Recipes 文档为主体,围绕日常任务队列开发中最常遇到的生产级问题展开:优雅关闭与中断任务重入队、队列监控、多队列拆分、Redis 高可用(Sentinel/Valkey/Cluster)、签名序列化、指数退避重试、任务去重、动态周期任务与 Web 框架(Flask/FastAPI)集成。全文将原文档的 16 个配方逐一展开,并结合仓库源码(huey/api.py、huey/signals.py、huey/serializer.py、huey/consumer_options.py 等)剖析底层实现,读完后你将获得一套可直接复制运行的 huey 生产级配置与编码方案。


1. 优雅关闭与中断任务重新入队(Graceful Shutdown & Re-enqueueing)

问题背景

当 consumer 被立即停止,或优雅关闭超出了--shutdown-timeout设定的等待时间时,任何正在执行中的任务都会被中断。默认情况下,这些任务会直接丢失(huey 的默认交付语义是 at-most-once)。如果任务不允许丢失,可以借助SIGNAL_INTERRUPTED信号把它们重新放回队列。

实现方式

通过huey.signal()注册一个监听SIGNAL_INTERRUPTED的处理函数,在回调里调用huey.enqueue(task)把任务重新入队:

from huey.signals import SIGNAL_INTERRUPTED @huey.signal(SIGNAL_INTERRUPTED) def on_interrupted(signal, task, *args, **kwargs): huey.enqueue(task)

从源码看,中断信号由 huey/api.py 的notify_interrupted_tasks()触发:当 consumer 停止时,它会将_tasks_in_flight中所有仍在执行的任务逐个取出并_emit(SIGNAL_INTERRUPTED, task)。同时,_execute()中捕获KeyboardInterrupt时也会发出SIGNAL_INTERRUPTED(见 huey/api.py)。信号的完整定义与分发机制见 huey/signals.py。

配套部署配置

要让大多数进程管理器发出的SIGTERM成为"优雅关闭"信号,需要以-g TERM启动 consumer,并将--shutdown-timeout设置在 supervisor 的 kill 截止时间之下几秒。

以supervisord为例,它默认会在stopwaitsecs(默认 10 秒)之后发送SIGKILL:

[program:my_huey] command=/path/to/venv/bin/huey_consumer my_app.huey -w 4 -g TERM -t 20 stopwaitsecs=30

这里-t 20给 worker 20 秒完成当前任务,supervisor 留出 30 秒等待窗口,二者留有安全余量。

相关 CLI 参数说明

从 huey/consumer_options.py 可以看到 consumer 的默认配置与参数:

参数短选项默认值说明
--workers-w1worker 线程/进程数量
--worker-type-kthread可选thread、greenlet、process
--shutdown-timeout-tNone(一直等)优雅关闭时等待 worker 完成任务的秒数,超时后中断任务
--graceful-signal-gINT触发优雅关闭的信号,可选INT或TERM;另一个信号则用于立即中断运行中的任务

完整的部署与信号讨论可参考 docs/deployment.rst(其中包含部署信号相关小节)。


2. 监控队列深度(Monitoring Queue Depth)

huey 提供了若干自省方法,非常适合构建监控端点或健康检查。核心是三个计数方法:

def get_queue_stats(): return { 'pending': huey.pending_count(), 'scheduled': huey.scheduled_count(), 'results': huey.result_count(), }

在 Web 框架中可以直接暴露为 HTTP 端点(以 Flask 为例):

# Flask example. @app.route('/huey/health/') def huey_health(): stats = get_queue_stats() return jsonify(stats)

这些方法在 huey/api.py 中实现:pending_count()调用存储层的queue_size(),scheduled_count()调用schedule_size(),result_count()调用result_store_size(),全部是 O(1) 量级的计数操作,不涉及反序列化。

如果需要查看队列中的实际任务,可以使用pending()/scheduled():

# 列出所有待处理任务(会反序列化每个任务,队列很大时可能较慢)。 for task in huey.pending(): print(task.name, task.id, task.args) # 列出所有已调度任务。 for task in huey.scheduled(): print(task.name, task.id, task.eta)

注意:pending()和scheduled()会反序列化队列中的每个任务(底层经由_deserialize_all(),见 huey/api.py),队列非常大时会很慢。当只需要计数时,请使用pending_count()和scheduled_count()。

作为补充,huey 自带的统计组件huey.contrib.stats的live_counts()正是用这三种计数方法实现的(见 huey/contrib/stats.py),它把所有存储调用包裹在try/except中,保证统计端点在后端暂时不可用时仍能返回None而非崩溃。


3. 用信号采集任务指标(Using Signals for Task Metrics)

信号机制同样可以用于为每个任务记录执行时长等指标:

import time from huey.signals import SIGNAL_EXECUTING, SIGNAL_COMPLETE, SIGNAL_ERROR _task_start_times = {} @huey.signal(SIGNAL_EXECUTING) def on_executing(signal, task): _task_start_times[task.id] = time.monotonic() @huey.signal(SIGNAL_COMPLETE, SIGNAL_ERROR) def on_finished(signal, task, exc=None): start = _task_start_times.pop(task.id, None) if start is not None: duration = time.monotonic() - start metrics.timing('huey.task.duration', duration, tags={ 'task': task.name, 'status': 'error' if signal == 'error' else 'ok', })

这里利用了huey.signal(*signals)可同时监听多个信号的能力(见 huey/api.py 与 huey/signals.py 的Signal.connect/send实现)。SIGNAL_EXECUTING在任务开始执行前发出,SIGNAL_COMPLETE/SIGNAL_ERROR在任务结束(成功或失败)时发出,分别对应 huey/api.py 与 huey/api.py。

两条重要的使用约束

  • 信号处理函数必须够快:信号处理器由 consumer 的 worker同步执行,一个慢的 handler 会阻塞 worker 拾取下一个任务。如果你的 metrics 客户端涉及网络 I/O,建议缓冲写入或使用异步客户端。
  • _task_start_times的内存模型:这是一个普通 dict。使用 thread 与 greenlet worker 时各 worker 共享同一进程内存,没问题;使用 process worker 时每个进程有自己的一份 dict,但因为给定任务只会在单个进程中执行,结果依然正确(对应-k三种 worker 类型的差异可参见 huey/constants.py)。

类似地,huey 自带的HueyStats组件也正是通过监听全量信号来生成事件记录:它在SIGNAL_EXECUTING时记录开始时间,在TERMINAL信号集合(complete/error/canceled/interrupted/expired/revoked/timeout/locked/rate-limited/retrying)上结算时长(见 huey/contrib/stats.py 与 huey/contrib/stats.py)。


4. 多队列(Multiple Queues)

何时需要第二个队列

一个 huey 应用由三部分组成:一个Huey实例(即一个具名队列)、注册到该实例上的任务、以及运行这些任务的 consumer 进程。

在引入第二个队列之前,先考虑任务优先级(priority)是否能解决你的问题(参见 docs/api.rst 中关于优先级的讨论)。使用独立队列的场景包括:

  • 针对不同类别的任务使用不同的并发度或 worker 类型,例如 CPU 密集任务用-k process,IO 密集任务用-k greenlet -w 50;
  • 隔离性:海量廉价任务绝不能挤占关键任务的 worker;
  • 按机器路由:某些任务只应在特定主机上运行。

声明多个实例

多个Huey实例可以共享同一个连接池:

# myapp/queues.py from huey import RedisHuey from redis import ConnectionPool pool = ConnectionPool(host='localhost', port=6379, max_connections=20) emails = RedisHuey('myapp_emails', connection_pool=pool) reports = RedisHuey('myapp_reports', connection_pool=pool)

任务按装饰器路由

任务归属由声明它所用的装饰器决定:

# myapp/tasks.py from huey import crontab from myapp.queues import emails, reports @emails.task(retries=2) def send_email(to, subject, body): ... @reports.task() def build_report(day): ... @reports.periodic_task(crontab(minute='0', hour='3')) def nightly_rollup(): ...

独立 consumer、独立调优

每个队列都有自己独立的 consumer,可以分别调参:

huey_consumer myapp.queues.emails -w 8 -k greenlet huey_consumer myapp.queues.reports -w 2 -k process

运行 N 个 consumer 意味着要监管 N 个进程。使用 systemd 时,可以用一个模板单元/etc/systemd/system/huey@.service,其中%i会展开为myapp.queues中的属性名:

[Unit] Description=huey consumer for %i [Service] WorkingDirectory=/srv/myapp ExecStart=/srv/myapp/venv/bin/huey_consumer myapp.queues.%i -w 4 -g TERM -t 55 TimeoutStopSec=60 Restart=on-failure [Install] WantedBy=multi-user.target
systemctl enable --now huey@emails huey@reports

多队列的存储层面

这一方案适用于任何存储后端。多个实例甚至可以共享同一个数据库文件:

emails = SqliteHuey('emails', filename='/var/lib/myapp/huey.db') reports = SqliteHuey('reports', filename='/var/lib/myapp/huey.db')

易踩的坑

  • 一切均按实例隔离:结果(results)、调度(schedules)、锁(locks)与撤销(revocations)都是 per-instance 的。一个Result句柄只能通过产生它的那个实例读取。
  • 周期任务归属:周期任务属于声明它的实例,并由该实例的 consumer 入队。-n/--no-periodic选项只有在运行同一队列的多个 consumer 时才需要(见 docs/consumer.rst 中 multiple consumers 小节)。
  • 队列名即存储命名空间:在相同存储上同名的两个实例就是同一个队列。Redis 存储会对名字做净化,剔除除字母数字和下划线之外的所有字符,因此队列名在净化之后仍必须保持互不相同。
  • immediate 模式是 per-instance 的设置:在测试中,记得对每个实例都打开它。

关于 supervisor 配置的完整讨论可参考 docs/deployment.rst。Django 用户如需以settings.HUEY风格配置多队列,可以查看第三方django-huey包。


5. 高可用 Redis:Sentinel、Valkey 与 Cluster

Sentinel 开箱即用

RedisHuey接受一个预先配置好的connection_pool,因此 huey 可以直接配合 Redis Sentinel 使用。master_for()返回的连接池在故障转移(failover)之后会自动重新发现当前 master:

from redis.sentinel import Sentinel from huey import RedisHuey sentinel = Sentinel( [('10.0.0.1', 26379), ('10.0.0.2', 26379), ('10.0.0.3', 26379)], # 重要:socket 超时必须明显大于 consumer 的阻塞读取超时(默认 1 秒)。 # 否则每次阻塞出队都会在 socket 层超时,恰好在那一刻弹出的任务 # 会被交付给一个已关闭的 socket 而丢失,且不会记录任何错误。 socket_timeout=5.0) huey = RedisHuey( 'my-app', connection_pool=sentinel.master_for('my-master').connection_pool)

必须使用 master_for()

永远使用master_for()。huey 在每条代码路径上都会读和写,所以使用slave_for()连接池会以ReadOnlyError失败。如果部署需要认证,为数据节点传入password=,为 sentinel 自身传入sentinel_kwargs={'password': '...'}。

Django 用户可以通过直接赋值实例来配置:

# settings.py from redis.sentinel import Sentinel from huey import RedisHuey sentinel = Sentinel([('10.0.0.1', 26379), ...], socket_timeout=5.0) HUEY = RedisHuey( 'my-app', connection_pool=sentinel.master_for('my-master').connection_pool)

故障转移(Failover)行为

consumer 在构造上就是容忍故障转移的:当 master 宕机时,worker 记录出队错误并应用退避,scheduler 记录错误并在下一个调度周期重试,一旦 Sentinel 提升新的 master,连接池会自动重连,无需重启。

但有两个事实必须理解:

  1. Redis 复制是异步的:落在故障转移窗口内的写入(入队、结果)可能在旧 master 未复制的数据被丢弃时丢失。
  2. 出队是破坏性的:已经被交付给 worker 的任务只存在于该 worker 的内存中,如果 worker 在任务执行中死亡,任务就没了。这是 huey 正常的 at-most-once 行为,并非 Sentinel 引入的问题。

因此,高可用 Redis 并不会让单个任务变得持久。请搭配幂等任务设计以及本文第 1 节的 中断任务重入队配方 一起使用。

Valkey 与 Redict

Valkey 和 Redict 与 Redis 在线协议兼容,可直接搭配RedisHuey使用:只需把地址指向 valkey 主机即可(redis-py客户端两者都支持)。如果更喜欢官方的valkey-glide客户端,huey 在huey.contrib.valkey_glide中提供了ValkeyGlideHuey(实现见 huey/contrib/valkey_glide.py,其底层存储继承自RedisStorage并适配了 glide 的同步 API)。

Redis Cluster

RedisHuey的 Redis 用法是单 key 的(唯一的一个 Lua 脚本自 2.4.2 起已兼容 cluster),但它构造的是标准的 redis-py 客户端与连接池,并不直接支持 redis-py 的RedisCluster客户端。因此对于高可用场景,Sentinel 才是被支持的拓扑。


6. 面向不受信环境的签名序列化器(Signed Serializer)

风险与对策

默认情况下 huey 使用pickle序列化任务与结果。如果 Redis 实例是共享的或暴露在网络中,恶意攻击者可能注入精心构造的 pickle 载荷。SignedSerializer为每条消息附加一个 HMAC 签名,任何被篡改的数据都会被拒绝:

from huey import RedisHuey from huey.serializer import SignedSerializer huey = RedisHuey( 'my-app', serializer=SignedSerializer(secret='my-secret-key'))

使用要点

  • secret在应用进程与 consumer 进程中必须是同一个值。
  • 如果消息被篡改,反序列化时会抛出ValueError。

从实现看(huey/serializer.py),SignedSerializer在序列化时先执行父类Serializer._serialize()(即pickle.dumps),再用hmac.new(key, message, hashlib.sha1)计算签名并追加message + ':' + signature;反序列化时先通过_unsign()用常数时间比较(hmac.compare_digest)校验签名,校验失败即抛ValueError。

注意:签名序列化器不加密数据,它只能检测篡改,任务参数在 Redis 中仍然是明文可见的。如果需要加密,可以继承Serializer并用自己的_serialize/_deserialize方法实现(比如借助cryptography库)。


7. 指数退避重试(Exponential Backoff Retries)

为什么需要退避

调用外部服务时指数退避很重要。没有它,一群 worker 按相同间隔同时重试,会形成"惊群效应"(thundering herd),把正在恢复的服务再次压垮。

配置方式

在任务装饰器上同时指定retry_backoff倍数、retries与retry_delay。首次重试等待retry_delay秒,之后每次延迟乘以retry_backoff:

@huey.task(retries=5, retry_delay=2, retry_backoff=2) def call_external_api(endpoint, payload): resp = requests.post(endpoint, json=payload) resp.raise_for_status() return resp.json()

重试时间表

如果 consumer 在12:00:00开始执行任务,重试计划如下:

  • 12:00:00首次调用
  • 12:00:02重试 1(延迟 2s)
  • 12:00:06重试 2(延迟 4s)
  • 12:00:14重试 3(延迟 8s)
  • 12:00:30重试 4(延迟 16s)
  • 12:01:02重试 5(延迟 32s)

从 huey/api.py 的_requeue_task()可以看到底层实现:每次重试入队前,task.retries -= 1,若配置了retry_backoff,则task.retry_delay *= task.retry_backoff,更新后的延迟随任务一起序列化,因此在多次重试之间可以持续增长。

用 expires 封顶

退避增长没有上限,所以较大的retries会把最后几次尝试推到几小时甚至几天之后。请用绝对expires日期时间来封顶整条重试链。相对expires(秒或 timedelta)会在每次重试重新入队时再次解析,因此它并不能限制整条链。

重试语义细节

  • 显式重试时间优先,且不会推进退避:RetryTask(delay=n)与RetryTask(eta=...)使用自己的延迟。
  • 被限速(rate-limited)的重试会等待当前延迟,而不是限速窗口,并且不推进退避。
  • 裸的RetryTask()遵循上述退避时间表。
  • 当前延迟随任务一起传递,所以一个已调度的任务报告的是它下一次尝试的等待时间,而不是声明的retry_delay。

相关异常类型定义见 huey/exceptions.py(RetryTask)与 huey/exceptions.py(重试类异常集合)。


8. 自定义错误元数据(Custom Error Metadata)

覆盖 build_error_result

任务失败时,huey 会存储一个错误结果。可以通过继承 Huey 实例并覆盖build_error_result来丰富它:

from huey import RedisHuey class MyHuey(RedisHuey): def build_error_result(self, task, exception): err = super().build_error_result(task, exception) err['task_name'] = task.name err['task_args'] = task.args err['task_kwargs'] = task.kwargs return err huey = MyHuey('my-app')

默认实现(huey/api.py)会生成包含error(异常 repr)、retries、traceback(完整 traceback 字符串)和task_id四个字段的字典;你的子类只需追加自定义字段。

捕获端读取自定义字段

现在当捕获TaskException时,metadata 字典就包含自定义字段了:

result = failing_task('some-arg') try: result.get(blocking=True, timeout=10) except TaskException as exc: print(exc.metadata['task_name']) # 'failing_task' print(exc.metadata['task_args']) # ('some-arg',) print(exc.metadata['traceback']) # 完整 traceback 字符串

从 huey/api.py 可以看到,Result.get()在取到Error类型的载荷时会抛出TaskException(result.metadata),而TaskException.metadata正是错误结果字典(见 huey/exceptions.py)。


9. 任务去重(Task Deduplication)

用 KV 存储 + pre_execute 钩子去重

可以利用键值存储和pre_execute钩子,在相同任务已在运行时跳过它:

import hashlib def dedup_key(task): digest = hashlib.md5(repr(task.data).encode()).hexdigest() return 'dedup:%s:%s' % (task.name, digest) @huey.pre_execute() def deduplicate(task): if not huey.put_if_empty(dedup_key(task), '1'): raise CancelExecution('Duplicate task, skipping') @huey.post_execute() def clear_dedup(task, task_value, exc): huey.delete(dedup_key(task))

huey.put_if_empty()是原子操作(见 huey/api.py,底层调用存储的put_if_empty,Redis 实现对应HSETNX),因此并发场景下只有一个 worker 能成功写入去重键,其余 worker 会抛出CancelExecution被跳过。

pre_execute/post_execute钩子的注册与执行机制见 huey/api.py(注册)与 huey/api.py(_run_pre_execute/_run_post_execute:前者在任务真正执行前运行,若抛CancelExecution则中止任务并发出SIGNAL_CANCELED)。注意CancelExecution定义于 huey/exceptions.py。


10. 键值数据存储(Key/Value Data Storage)

huey 的 result-store 本身可以直接当作一个便利的任意键值缓存来用:

@huey.task() def calculate_something(): # 默认情况下 result store 把 get() 当作 pop() 处理, # 为了保留数据以便再次读取,需要传入第二个参数 peek=True。 prev_results = huey.get('calculate-something.result', peek=True) if prev_results is None: # 没有历史结果,从头开始计算。 data = start_from_beginning() else: # 只计算自上次以来变化的部分。 data = just_what_changed(prev_results) # 把更新后的数据存回 result store。 huey.put('calculate-something.result', data) return data

底层方法见 huey/api.py:put()序列化后调用存储的put_data();get()默认走pop_data()(读取即删除),peek=True时走peek_data()(只读不删);delete()调用delete_data()。存储层接口定义于 huey/storage.py。Huey.get/Huey.put的更多细节可参考 docs/api.rst。


11. 用键值存储做进度跟踪(Progress Tracking via Key/Value Storage)

长任务往往需要向调用方汇报进度。huey 的键值存储(Huey.put/Huey.get)是实现这一点的便捷途径:

@huey.task(context=True) def process_large_file(filepath, task=None): lines = open(filepath).readlines() total = len(lines) results = [] for i, line in enumerate(lines): results.append(transform(line)) if i % 100 == 0: huey.put('progress:%s' % task.id, { 'current': i, 'total': total, 'pct': int(100 * i / total), }) huey.put('progress:%s' % task.id, { 'current': total, 'total': total, 'pct': 100, }) return results

这里使用了context=True:任务执行时会把任务实例注入kwargs['task'](实现见 huey/api.py 的execute定义),从而在任务内部拿到task.id。

调用方可以轮询进度:

result = process_large_file('/data/big.csv') # 轮询进度。使用 peek=True 以读取而不删除。 progress = huey.get('progress:%s' % result.id, peek=True) if progress: print('%d%% complete' % progress['pct'])

注意:默认Huey.get是破坏性的(读取后删除该值)。传入peek=True可只读不删。用完记得清理进度键:huey.delete('progress:%s' % task_id)。


12. 动态扇出(Dynamic Fan-Out)

Chord 的成员必须在入队时就确定。当子任务集合依赖运行时的值(例如分页的 API 结果)时,可以在任务内部发起 chord 入队:

@huey.task() def fetch_page(url): return requests.get(url).json() @huey.task() def aggregate(results): combined = {} for page_data in results: combined.update(page_data) return combined @huey.task() def discover_and_fetch(base_url): index = requests.get(base_url).json() urls = [item['url'] for item in index['items']] result = huey.enqueue( chord([fetch_page.s(u) for u in urls], aggregate.s())) # 可选:保存回调任务的 ID,方便调用方追踪。 huey.put('fanout-result-id', result.callback.id)

这里的chord([...], callback)会在所有成员任务完成后自动入队回调任务aggregate(并聚合各成员的返回值列表)。底层实现见 huey/api.py 的_enqueue_chord():它为每个成员构造ChordConfig,最后一个完成的成员触发_check_chord(),在所有结果就绪后把(results,)作为参数喂给回调并enqueue(callback)(见 huey/api.py)。


13. 动态周期任务(Dynamic Periodic Tasks)

原理与限制

要动态创建周期任务,必须把它注册到由 consumer 的 scheduler 线程所维护的内存调度表中。由于这个注册表在内存中,任何动态定义的任务都必须在最终执行调度的进程——即 consumer——内注册。

警告:以下示例在processworker 类型下不生效,因为目前没有办法与调度进程交互。使用 thread 或 greenlet 时,worker 线程与 scheduler 线程共享同一份内存调度表,因此可以进行修改。

实现

def dynamic_ptask(message): print('dynamically-created periodic task: "%s"' % message) @huey.task() def schedule_message(message, cron_minutes, cron_hours='*'): def wrapper(): dynamic_ptask(message) schedule = crontab(cron_minutes, cron_hours) # 需要为任务提供唯一名字。方法有很多(基于参数等), # 这里简单用 uuid 即可。 task_name = 'dynamic_ptask_%s' % uuid.uuid4().hex huey.periodic_task(schedule, name=task_name)(wrapper)

注意huey.periodic_task()在 huey/api.py 中的签名是periodic_task(validate_datetime, retries=0, retry_delay=0, retry_backoff=0, ...),第一个参数是一个校验函数(crontab返回的可调用对象),通过装饰器把它包装为PeriodicTask并注册进实例的 registry。

假设 consumer 正在运行,现在可以创建任意多个"动态周期任务"实例:

>>> from demo import schedule_message >>> schedule_message('I run every 5 minutes', '*/5') <Result: task ...> >>> schedule_message('I run between 0-15 and 30-45', '0-15,30-45') <Result: task ...>

任务注册后,scheduler 会在每次调度周期检查read_periodic()(见 huey/api.py),把所有validate_datetime(timestamp)返回 True 的周期任务入队。


14. 把任意函数当作任务运行(Run Arbitrary Functions as Tasks)

与其预先显式声明所有任务,不如写一个通用任务,它接收一个点分导入路径并调用任意函数:

from importlib import import_module @huey.task() def path_task(path, *args, **kwargs): module_path, name = path.rsplit('.', 1) mod = import_module(module_path) return getattr(mod, name)(*args, **kwargs) # 用法:在 consumer 进程中运行 myapp.utils.reindex('products')。 path_task('myapp.utils.reindex', 'products')

警告:谨慎使用此模式。被调用的函数必须能被 consumer 进程 import,参数必须可 pickle。由于它可以调用任何可导入的可调用对象,不要把它暴露给不受信任的输入。

这个配方同样绕过了 huey 在 huey/api.py 中"不支持 async 函数"的限制——普通def函数在这里没有此约束,但path_task调用的目标仍必须是普通同步函数。


15. 在 Flask 中使用 Huey

huey 与框架无关,不需要任何扩展即可配合 Flask 使用。

声明实例与任务

与应用同时(或之前)声明实例,并从任务模块导入它:

# app.py from flask import Flask from huey import RedisHuey app = Flask(__name__) huey = RedisHuey('my-app')
# tasks.py from app import app, huey @huey.task() def send_welcome_email(user_id): # 任务运行在 consumer 进程中,不在任何请求内。 # 如果任务使用 Flask 扩展(数据库、邮件等),需要提供应用上下文: with app.app_context(): user = User.query.get(user_id) mail.send(make_welcome_message(user))

视图入队

# views.py from app import app from tasks import send_welcome_email @app.route('/signup/', methods=['POST']) def signup(): user = create_user(request.form) send_welcome_email(user.id) return redirect(url_for('welcome'))

用装饰器收敛样板代码

如果大多数任务都需要应用上下文,可以把样板封装成一个小装饰器:

import functools def flask_task(*task_args, **task_kwargs): def decorator(fn): @functools.wraps(fn) def inner(*args, **kwargs): with app.app_context(): return fn(*args, **kwargs) return huey.task(*task_args, **task_kwargs)(inner) return decorator @flask_task(retries=2) def send_welcome_email(user_id): ...

启动 consumer

consumer 指向一个导入 app 与全部任务的入口模块(参见 docs/imports.rst):

# main.py from app import app, huey import tasks
huey_consumer main.huey -w 4

一个完整可运行的示例应用位于 examples/flask_ex,其中包含了 examples/flask_ex/app.py、examples/flask_ex/tasks.py、examples/flask_ex/views.py 以及启动脚本 examples/flask_ex/run_huey.sh 与 examples/flask_ex/run_webapp.sh。


16. 在 FastAPI 中使用 Huey

独立模块声明

在独立模块中声明实例与任务:

# tasks.py from huey import RedisHuey huey = RedisHuey('my-app') @huey.task() def generate_report(user_id): ... # 重活都在 consumer 里干。 return report_data

异步请求中入队与轮询

从异步请求处理器中入队完全没问题,因为那只是一次快速的单次存储写入:

# api.py from fastapi import FastAPI from huey.contrib.asyncio import aget_result from tasks import generate_report, huey app = FastAPI() @app.post('/report/{user_id}') async def begin_report(user_id: int): rh = generate_report(user_id) return {'task_id': rh.id} @app.get('/report/status/{task_id}') async def report_status(task_id: str): # 对 result store 的非阻塞、非破坏性读取。 value = huey.result(task_id, preserve=True) return {'ready': value is not None, 'value': value}

阻塞等待任务结果

要一直保持请求打开直到任务完成,可以配合 docs/asyncio.rst 中描述的 asyncio 辅助函数等待结果,期间其他请求仍可继续被服务:

@app.post('/report/{user_id}/wait') async def report_wait(user_id: int): rh = generate_report(user_id) value = await aget_result(rh) return {'value': value}

aget_result位于 huey/contrib/asyncio.py,对应的测试用例见 huey/tests/test_asyncio.py。

要点与注意事项

  • consumer 单独运行:huey_consumer tasks.huey -w 4。
  • 任务函数是普通def函数,huey不会执行async def任务(见 huey/api.py 的显式检查)。IO 密集负载仍可通过-k greenlet在 consumer 中获得高并发。
  • 如果任务抛出了异常,读取其结果会抛出TaskException,如果任务可能失败,请在状态端点中处理它。
  • huey.result(task_id)默认是破坏性的;preserve=True会保留结果,以便状态端点被反复轮询。

结语

以上 16 个配方覆盖了 huey 在生产环境中最高频的真实需求:从保证任务不丢的优雅关闭与中断重入队,到可观测的队列监控与信号指标;从多队列隔离、Redis 高可用,到安全序列化、退避重试、去重、进度追踪乃至与主流 Web 框架的集成。每个配方都在原文档基础上补充了 huey/api.py、huey/signals.py、huey/serializer.py、huey/exceptions.py、huey/consumer_options.py、huey/storage.py 等处的源码级佐证,你可以放心地把它们组合进自己的应用。完整的 API 参考与更多机制说明,可继续阅读 docs/api.rst、docs/consumer.rst 与 docs/deployment.rst。

  • 任务调度
  • 后端

【免费下载链接】huey

a little task queue for python

项目地址:https://gitcode.com/gh_mirrors/hu/huey
点击查看免费下载

相关推荐

上一篇:如何通过代理抓包技术实现跨平台网络资源下载
下一篇:ESP-IDF 命令行前端工具 idf.py 完全指南:项目构建、烧录、调试与扩展

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询