资讯详情

资讯详情

Huey 实战配方手册:优雅关闭、多队列、Redis 高可用与 16 个任务队列高级用法

任务调度后端【免费下载链接】hueya 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 stopwaitsecs30这里-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 Depthhuey 提供了若干自省方法非常适合构建监控端点或健康检查。核心是三个计数方法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, excNone): 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 processIO 密集任务用-k greenlet -w 50隔离性海量廉价任务绝不能挤占关键任务的 worker按机器路由某些任务只应在特定主机上运行。声明多个实例多个Huey实例可以共享同一个连接池# myapp/queues.py from huey import RedisHuey from redis import ConnectionPool pool ConnectionPool(hostlocalhost, port6379, max_connections20) emails RedisHuey(myapp_emails, connection_poolpool) reports RedisHuey(myapp_reports, connection_poolpool)任务按装饰器路由任务归属由声明它所用的装饰器决定# myapp/tasks.py from huey import crontab from myapp.queues import emails, reports emails.task(retries2) def send_email(to, subject, body): ... reports.task() def build_report(day): ... reports.periodic_task(crontab(minute0, hour3)) 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] Descriptionhuey consumer for %i [Service] WorkingDirectory/srv/myapp ExecStart/srv/myapp/venv/bin/huey_consumer myapp.queues.%i -w 4 -g TERM -t 55 TimeoutStopSec60 Restarton-failure [Install] WantedBymulti-user.targetsystemctl enable --now hueyemails hueyreports多队列的存储层面这一方案适用于任何存储后端。多个实例甚至可以共享同一个数据库文件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. 高可用 RedisSentinel、Valkey 与 ClusterSentinel 开箱即用RedisHuey接受一个预先配置好的connection_pool因此 huey 可以直接配合 Redis Sentinel 使用。master_for()返回的连接池在故障转移failover之后会自动重新发现当前 masterfrom 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_timeout5.0) huey RedisHuey( my-app, connection_poolsentinel.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_timeout5.0) HUEY RedisHuey( my-app, connection_poolsentinel.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 与 RedictValkey 和 Redict 与 Redis 在线协议兼容可直接搭配RedisHuey使用只需把地址指向 valkey 主机即可redis-py客户端两者都支持。如果更喜欢官方的valkey-glide客户端huey 在huey.contrib.valkey_glide中提供了ValkeyGlideHuey实现见 huey/contrib/valkey_glide.py其底层存储继承自RedisStorage并适配了 glide 的同步 API。Redis ClusterRedisHuey的 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, serializerSignedSerializer(secretmy-secret-key))使用要点secret在应用进程与 consumer 进程中必须是同一个值。如果消息被篡改反序列化时会抛出ValueError。从实现看huey/serializer.pySignedSerializer在序列化时先执行父类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_backoffhuey.task(retries5, retry_delay2, retry_backoff2) def call_external_api(endpoint, payload): resp requests.post(endpoint, jsonpayload) resp.raise_for_status() return resp.json()重试时间表如果 consumer 在12:00:00开始执行任务重试计划如下12:00:00首次调用12:00:02重试 1延迟 2s12:00:06重试 2延迟 4s12:00:14重试 3延迟 8s12:00:30重试 4延迟 16s12: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(delayn)与RetryTask(eta...)使用自己的延迟。被限速rate-limited的重试会等待当前延迟而不是限速窗口并且不推进退避。裸的RetryTask()遵循上述退避时间表。当前延迟随任务一起传递所以一个已调度的任务报告的是它下一次尝试的等待时间而不是声明的retry_delay。相关异常类型定义见 huey/exceptions.pyRetryTask与 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(blockingTrue, timeout10) 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_emptyRedis 实现对应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 Storagehuey 的 result-store 本身可以直接当作一个便利的任意键值缓存来用huey.task() def calculate_something(): # 默认情况下 result store 把 get() 当作 pop() 处理 # 为了保留数据以便再次读取需要传入第二个参数 peekTrue。 prev_results huey.get(calculate-something.result, peekTrue) 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.pyput()序列化后调用存储的put_data()get()默认走pop_data()读取即删除peekTrue时走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(contextTrue) def process_large_file(filepath, taskNone): 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这里使用了contextTrue任务执行时会把任务实例注入kwargs[task]实现见 huey/api.py 的execute定义从而在任务内部拿到task.id。调用方可以轮询进度result process_large_file(/data/big.csv) # 轮询进度。使用 peekTrue 以读取而不删除。 progress huey.get(progress:%s % result.id, peekTrue) if progress: print(%d%% complete % progress[pct])注意默认Huey.get是破坏性的读取后删除该值。传入peekTrue可只读不删。用完记得清理进度键huey.delete(progress:%s % task_id)。12. 动态扇出Dynamic Fan-OutChord 的成员必须在入队时就确定。当子任务集合依赖运行时的值例如分页的 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, nametask_name)(wrapper)注意huey.periodic_task()在 huey/api.py 中的签名是periodic_task(validate_datetime, retries0, retry_delay0, retry_backoff0, ...)第一个参数是一个校验函数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 中使用 Hueyhuey 与框架无关不需要任何扩展即可配合 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(retries2) def send_welcome_email(user_id): ...启动 consumerconsumer 指向一个导入 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, preserveTrue) 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)默认是破坏性的preserveTrue会保留结果以便状态端点被反复轮询。结语以上 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。赞分享任务调度后端【免费下载链接】hueya little task queue for python项目地址https://gitcode.com/gh_mirrors/hu/huey点击查看免费下载相关推荐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 任务队列瘫痪是否在寻找一种简单可靠的方案来保障任务处理的任务调度后端消息队列上一篇如何通过代理抓包技术实现跨平台网络资源下载下一篇ESP-IDF 命令行前端工具 idf.py 完全指南项目构建、烧录、调试与扩展创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →