
Celery 信号机制源码解析celery.utils.dispatch.signal 的 Observer 模式实现与实战【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celeryCelery 的分布式任务体系中有大量事件任务发布、任务执行成功/失败、Worker 启动/关闭、Beat 心跳等需要对外广播而承载这套事件通知机制的正是celery.utils.dispatch.signal模块中的Signal类。本文以该模块的 API 参考文档为主体结合源码实现与真实调用方深入讲解Signal的完整 API、弱引用与线程安全设计、retry 重试机制并给出可直接复用的信号连接与调试实战方案。模块定位Celery 的 Observer 模式实现celery.utils.dispatch.signal是 Celery 内部对观察者模式Observer pattern的标准实现。模块通过 celery/utils/dispatch/init.py 对外仅暴露一个核心类Observer pattern. from .signal import Signal __all__ (Signal,)从仓库历史看该实现源于 Django 的django.dispatch而django.dispatch又 fork 自 PyDispatcher许可声明保存在 celery/utils/dispatch/license.txt 中。因此它天然兼容 Django 信号的惯用写法但又在弱引用、线程安全和重试等方面做了 Celery 特有的增强。模块内公开的完整 API 由 celery/utils/dispatch/signal.py 提供包括构造器Signal(providing_argsNone, use_cachingFalse, nameNone)连接与断开connect()、disconnect()发送与查询send()、send_robust、has_listeners()内部机制_live_receivers()、_clear_dead_receivers()、_remove_receiver()等Signal 构造器与内部状态Signal的__init__接受三个可选参数见 signal.py参数类型默认值说明providing_argsList / Set[]该信号在send()时会携带的参数名清单主要用于文档与调试不强制校验use_cachingboolFalse是否按 sender 缓存接收者列表可显著减少重复过滤开销namestrNone信号名称用于调试体现在__repr__输出中构造时同时初始化以下关键内部状态self.receivers []接收者注册表元素为(lookup_key, receiver)二元组lookup_key由接收者与 sender 共同决定。self.lock threading.Lock()连接、断开、清理死引用时使用的互斥锁保证多线程环境下的安全。self.sender_receivers_cache当use_cachingTrue时是一个weakref.WeakKeyDictionary用于按 sender 缓存已解析的活跃接收者否则为空字典。self._dead_receivers False标记当前注册表中是否存在已失效的弱引用供惰性清理使用。__repr__会输出信号名与providing_args便于日志排查def __repr__(self): return f{type(self).__name__}: {self.name} providing_args{self.providing_args!r}connect()注册接收者的完整语义connect()是使用频率最高的 API支持两种调用风格见 signal.py# 风格一直接作为装饰器 signal.connect def handler(senderNone, **kwargs): ... # 风格二显式传参 signal.connect(handler, senderobj, weakTrue, dispatch_uidmy-id)其参数语义如下参数说明receiver接收信号的函数或实例方法必须是可调用callable且可哈希的对象sender接收者关心的发送者传特定对象则只接收该 sender 的信号传None则接收任意 sender 的信号weak是否对接收者使用弱引用默认True为False时持有强引用dispatch_uid接收者的唯一标识符任意可哈希对象用于避免同一接收者被重复注册retry接收者抛出异常如ConnectionError时是否自动重试直至成功开启后强制使用强引用忽略weak连接阶段会依次执行以下检查与处理见 signal.py可调用性断言assert callable(receiver)非可调用对象直接报错。关键字参数校验调用fun_accepts_kwargs(receiver)实现于 celery/utils/functional.py检查接收者是否接受关键字参数否则抛出ValueError(Signal receiver must accept keyword arguments.)。这正是官方文档强调信号处理函数最好写成**kwargs的底层原因——信号发送时会注入signal、sender等额外关键字。PromiseProxy sender 支持如果sender是PromiseProxyCelery 的惰性代理对象定义于 celery/local.py则通过sender.__then__(self._connect_proxy, ...)在代理解析完成后延迟连接见_connect_proxy()。生成查找键_make_lookup_key(receiver, sender, dispatch_uid)决定接收者的唯一性。若显式提供dispatch_uid则以(dispatch_uid, id(sender))为键否则以id(receiver), id(sender)为键。弱引用包装_boundmethod_safe_weakref()会特殊处理绑定方法——普通对象用weakref.ref绑定方法用weakref.WeakMethod并返回其宿主实例__self__随后通过weakref.finalize(宿主实例, self._remove_receiver)注册析构回调实例被回收时自动将_dead_receivers置位。去重与登记在self.lock保护下先清理死引用再遍历self.receivers检查lookup_key是否已存在不存在才追加。任何变更后都会清空sender_receivers_cache保证缓存一致性。retry 重试包装的细节开启retryTrue时见 signal.py接收者会被包装为_try_receiver_over_time其内部调用 kombu 的retry_over_time()按指数退避策略重试并在每次失败时记录形如RECEIVER_RETRY_ERROR的日志其中humanize_seconds来自 celery/utils/time.py。这里有两个容易踩坑的设计未提供dispatch_uid时会自动以原函数 id 作为dispatch_uid并写入包装函数属性fun._dispatch_uid。这样后续以原函数为键也能正确查找_make_lookup_key中hasattr(receiver, _dispatch_uid)分支正是为此对应 Issue #9119。weak会被强制设为False因为重试包装必须持有接收者的强引用才能反复调用。disconnect()断开连接disconnect(receiverNone, senderNone, weakNone, dispatch_uidNone)见 signal.py从注册表中移除匹配的接收者返回True/False表示是否确实断开了某个接收者。注意事项使用弱引用时通常无需手动断开——宿主实例被 GC 后会自动清理。weak参数已被废弃传入会触发CDeprecationWarning警告类定义于 celery/exceptions.py。与connect一样断开后也会清空 sender 缓存。send() / send_robust信号派发send(sender, **named)见 signal.py将信号从 sender 派发给所有已连接的接收者返回[(receiver, response), ...]列表。核心流程短路优化若注册表为空或缓存中标记了NO_RECEIVERS直接返回空列表避免无谓加锁。解析活跃接收者_live_receivers()在锁内完成死引用清理、按 sender 过滤senderkey NONE_ID或senderkey senderkey、解引用弱引用等操作若开启缓存则写入/读取sender_receivers_cache无接收者时缓存NO_RECEIVERS哨兵值NO_RECEIVERS object()见 signal.py。逐个调用每个接收者以receiver(signalself, sendersender, **named)形式调用——signal与sender是自动注入的固定关键字。异常收集而非抛出接收者抛出的任何异常都会被logger.exception记录并以(receiver, exc)形式放进响应列表不会中断后续接收者。因此 Celery 的send与 Django 的send_robust行为一致send_robust send这一别名纯粹是为了兼容 Django 接口而保留源码注释明确说明了这一点。这个异常即返回值的设计对任务系统尤为重要信号接收者如监控插件出错时绝不应影响任务消息的发布与消费主链路。死引用清理机制_remove_receiver()只是把_dead_receivers置True见 signal.py真正的清理发生在持有锁的connect、disconnect、_live_receivers调用链中。注释解释了原因_remove_receiver是 GC 副作用回调可能在持锁状态下触发直接在此处操作锁内数据结构会造成死锁。_make_id()还处理了 Celery 特有的Proxy对象见 signal.py对Proxy先取_get_current_object()对bytes/str直接返回原值对应 Issue #2475 的修复对绑定方法取id(__func__)其余取id(obj)。has_listeners()快速探测has_listeners(senderNone)返回当前是否存在活跃接收者底层即bool(self._live_receivers(sender))。它常被用于仅在有人监听时才做昂贵计算的优化场景。在 Celery 中的真实应用Signal并非仅供内部使用的玩具而是支撑整个 Celery 信号体系的基石。1. 内置信号全集celery/signals.pycelery/signals.py 集中定义了 Celery 向用户暴露的全部信号均基于Signal构造并明确标注providing_args。按类别可分为任务生命周期信号before_task_publish发布前含body/exchange/routing_key/headers/properties/declare/retry_policy、after_task_publish、task_received、task_prerun、task_postrun、task_success、task_retry、task_failure、task_internal_error、task_revoked、task_rejected、task_unknown以及已废弃的task_sent源码注明 6.0 移除改用after_task_publish。Worker 生命周期信号celeryd_init、celeryd_after_setup、import_modules、worker_init、worker_before_create_process、worker_process_init、worker_process_shutdown、worker_ready、worker_shutdown、worker_shutting_down、heartbeat_sent。日志与运行时信号setup_logging、after_setup_logger、after_setup_task_logger、beat_init、beat_embedded_init、eventlet_pool_started/preshutdown/postshutdown/apply、user_preload_options。2. App 级信号与事件快照celery/app/base.py 中Celery应用对象内部使用Signal实现on_configure、on_after_configure、on_after_finalize、on_after_fork等应用钩子见 base.py并提供了app.signals之外可编程注册的入口。celery/events/snapshot.py 使用信号机制实现事件快照的周期性落库。celery/apps/beat.py 与 celery/contrib/testing/worker.py 也分别用到了本模块的Signal。3. 官方使用示例docs/userguide/signals.rst 给出了最典型的连接写法——用after_task_publish观察任务发布from celery.signals import after_task_publish after_task_publish.connect def task_sent_handler(senderNone, headersNone, bodyNone, **kwargs): # 协议 v2 下任务信息位于 headers 中 info headers if task in headers else body print(after_task_publish for task id {info[id]}.format(infoinfo))由于after_task_publish以任务名为 sender还可以按任务名精确过滤after_task_publish.connect(senderproj.tasks.add) def task_sent_handler(senderNone, headersNone, bodyNone, **kwargs): info headers if task in headers else body print(after_task_publish for task id {info[id]}.format(infoinfo))4. 集成测试佐证仓库的冒烟测试 t/smoke/tests/test_signals.py 覆盖了信号在实际 worker 环境下的收发链路可作为验证自定义信号连接是否正确的最小参考。实战要点与性能建议结合源码与官方文档使用Signal时有几点值得注意处理函数务必带**kwargsconnect会强制校验这一点且未来 Celery 新增参数时只有**kwargs风格才不会被破坏。弱引用与生命周期默认weakTrue若接收者是局部函数或临时对象可能因被 GC 而静默失效需要长期存在的监听器应持有强引用weakFalse或保证宿主实例存活。异常不会阻断主流程接收者异常会被收集进响应列表并记录日志主流程继续执行因此可放心在信号里做监控、埋点等非关键逻辑。缓存开关存在大量 sender 且接收者集合稳定的场景可开启use_cachingTrue缓存会随connect/disconnect自动失效无需手动维护。多线程安全connect/disconnect/_live_receivers均在threading.Lock保护下操作可安全用于多线程应用。小结celery.utils.dispatch.signal用约 350 行代码实现了一个生产级的事件分发内核兼容 Django 的 API 形态、基于弱引用的自动生命周期管理、线程安全注册表、sender 级过滤与缓存、异常隔离以及可选的重试包装。理解它的设计不仅能更安全地使用 Celery 的全部内置信号也能在需要时基于同一模式构建自己的解耦事件系统。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考