- 任务调度
- 后端
【免费下载链接】huey
a little task queue for python
本指南以 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 | -w | 1 | worker 线程/进程数量 |
--worker-type | -k | thread | 可选thread、greenlet、process |
--shutdown-timeout | -t | None(一直等) | 优雅关闭时等待 worker 完成任务的秒数,超时后中断任务 |
--graceful-signal | -g | INT | 触发优雅关闭的信号,可选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.targetsystemctl 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,连接池会自动重连,无需重启。
但有两个事实必须理解:
- Redis 复制是异步的:落在故障转移窗口内的写入(入队、结果)可能在旧 master 未复制的数据被丢弃时丢失。
- 出队是破坏性的:已经被交付给 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 taskshuey_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
相关推荐
gocsv最佳实践总结:从新手到专家的10个关键技巧
gocsv最佳实践总结:从新手到专家的10个关键技巧 gocsv是Go语言中一款强大的CSV序列化与反序列化工具,它提供了简洁易用的API,帮助开发者轻松处理C
开发工具Langflow 多 worker 高可用部署怎么配置 Redis 任务队列?
Langflow 多 worker 高可用部署怎么配置 Redis 任务队列? Langflow 默认只运行一个 worker 进程,build 任务的状态(j
人工智能大模型AI AgentRAG后端前端MCP 服务工作流自动化Kue Redis 哨兵配置:实现高可用的任务队列服务
Kue Redis 哨兵配置:实现高可用的任务队列服务 你是否曾因 Redis 单点故障导致 Kue 任务队列瘫痪?是否在寻找一种简单可靠的方案来保障任务处理的
任务调度后端消息队列
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考