资讯动态

Celery Eventlet 并发池实战:从 urlopen 到递归爬虫与批量任务生产者的完整示例解析

发布时间:2026/9/20 14:09:11 来源:尧图企业网站定制
Celery Eventlet 并发池实战从 urlopen 到递归爬虫与批量任务生产者的完整示例解析【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery导读Eventlet 是 Celery 官方支持的异步并发执行池之一它通过协程greenlet让开发者沿用熟悉的阻塞式编程风格却能获得非阻塞 I/O 带来的高并发能力。本文以仓库 examples/eventlet 目录下的完整示例应用为骨架依次讲解 Eventlet 池的安装与 worker 启动方式、三个典型示例任务HTTP 探测、递归网页爬虫、批量任务生产者的用法与源码实现并深入到celery/concurrency/eventlet.py的 TaskPool 底层帮助你掌握在 I/O 密集场景下用 Celery Eventlet 写出可运行、可扩展的并发任务代码。示例应用概览examples/eventlet 目录结构该目录是一个独立的 Celery 应用包含一个配置文件与三个可运行示例文件作用celeryconfig.py应用配置broker 地址、导入模块、结果过期时间等tasks.py示例任务urlopen抓取 URL 并返回响应体字节数webcrawler.py示例任务crawl限定同域名的递归网页爬虫bulk_task_producer.pyProducerPool单进程内批量、并发发布任务的工具类注意原 README 在 Introduction 中写的是 two example tasks实际上该目录提供了三个独立示例urlopen、crawl、ProducerPool下文逐一展开。环境准备安装 Eventlet 与依赖根据 examples/eventlet/README.rst首先需要安装 Eventlet同时推荐安装dnspython——当它存在时Celery 的所有名称查找DNS 解析都会变为异步操作这对高并发网络任务非常关键$ python -m pip install eventlet celery pybloom-live其中eventlet并发库本体仓库的 requirements/extras/eventlet.txt 中声明的最低版本为eventlet0.32.0且仅针对python_version3.10的约束实际安装时建议使用满足项目要求的最新版本pybloom-live供webcrawler.py使用实现布隆过滤器BloomFilter用于降低重复爬取同一 URL 的概率dnspython可选但强烈推荐让 DNS 名称查找也走异步路径。除 Python 依赖外示例默认通过 RabbitMQ 作为 brokerceleryconfig.py 中broker_url amqp://guest:guestlocalhost:5672//因此需要先保证 RabbitMQ 已在localhost:5672运行默认账号密码为guest/guest。若尚未安装可参考 docs/getting-started/first-steps-with-celery.rst 入门指南。启动 Eventlet worker--pool 与高并发参数安装完成后进入示例目录并启动 worker$ cd examples/eventlet $ celery worker -l INFO --concurrency500 --pooleventlet关键点解读--pooleventlet等价于-P eventlet显式指定执行池为 Eventlet。这一点至关重要——celeryconfig.py 的注释明确提醒Never use the worker_pool setting as thatll patch the worker too late即不要通过配置文件里的worker_pool设置来启用 Eventlet因为那样会让 monkey patch 时机过晚导致补丁失效。必须通过命令行参数在进程启动早期完成补丁。--concurrency500等价于-c 500将并发度设为 500。这正是 Eventlet 的核心优势——docs/userguide/concurrency/eventlet.rst 指出prefork 池受限于每 CPU 几个进程而 Eventlet 可以高效地派生数百甚至上千个绿色线程green thread非常适合大量并发的网络 I/O 任务。-l INFO日志级别。配置文件中的其他设置broker_url amqp://guest:guestlocalhost:5672// worker_disable_rate_limits True result_expires 30 * 60 imports (tasks, webcrawler)worker_disable_rate_limits关闭限速避免在高并发下被速率限制拖慢result_expires结果过期时间为 30 分钟importsworker 启动时自动导入tasks与webcrawler两个模块从而注册其中的shared_task另外该文件还把当前工作目录插入sys.pathsys.path.insert(0, os.getcwd())保证tasks等模块可以直接按名导入。示例一tasks.urlopen——最简单的 I/O 并发任务任务源码tasks.py 的实现非常简洁import requests from celery import shared_task shared_task() def urlopen(url): print(f-open: {url}) try: response requests.get(url, timeout10.0) except requests.exceptions.RequestException as exc: print(f-url {url} gave error: {exc!r}) return return len(response.text)它接收一个 URL用requests.get发起 HTTP 请求带 10 秒超时成功则返回响应正文的字符数失败则打印异常并返回None。请求异常被完整捕获因此单个 URL 失败不会影响其他任务的执行。单任务调用README 中的交互式调用方式$ cd examples/eventlet $ python from tasks import urlopen urlopen.delay(https://www.google.com/).get() 9980delay()把任务异步发送到 broker.get()阻塞等待结果返回的9980是当时 Google 首页响应正文的字符数该数值随页面内容变化这里仅为示例输出。批量并发调用要同时打开多个 URL可以借助 Celery 的 group 原语构造并行签名组 from tasks import urlopen from celery import group result group(urlopen.s(url) ... for url in LIST_OF_URLS).apply_async() for incoming_result in result.iter_native(): ... print(incoming_result)这里urlopen.s(url)创建任务的签名signaturegroup(...)把它们组合成一个并行组一次性提交。iter_native()是 celery/result.py 中定义的后端优化版迭代器逐个产出已完成的结果而无需等待全部完成其文档字符串明确说明该能力目前仅由amqp、Redis 和 cache三类结果后端支持使用前请确认后端类型。示例二webcrawler.crawl——受限域名的递归爬虫设计思路webcrawler.py 是一个递归爬虫示例其模块文档字符串说明了几项关键设计决策只爬取**当前主机名域名**下的 URL避免摧毁互联网代码注释原话使用pybloom_live.BloomFilter布隆过滤器记录已见过的 URL以较低的误判率换取极小的内存占用布隆过滤器不作为共享服务而是作为参数传递给每个子任务——文档也承认更优做法是把去重放到集中式服务例如 Redis 的 set中任务设置serializerpickle、compressionzlib容量 100,000 成员、误判率 0.001 的 BloomFilter 在 pickle 后约 2.8MB但经 zlib 压缩后仅约 2.9kB因此直接在任务层面启用压缩即可无需手动处理。核心实现shared_task(ignore_resultTrue, serializerpickle, compressionzlib) def crawl(url, seenNone): print(fcrawling: {url}) if not seen: seen BloomFilter(capacity50000, error_rate0.0001) with Timeout(5, False): try: response requests.get(url) except requests.exception.RequestError: return location domain(url) wanted_urls [] for url_match in url_regex.finditer(response.text): url url_match.group(0) # To not destroy the internet, we only fetch URLs on the same domain. if url not in seen and location in domain(url): wanted_urls.append(url) seen.add(url) subtasks group(crawl.s(url, seen) for url in wanted_urls) subtasks.delay()逐段解读ignore_resultTrue爬虫只关心扩散不收集返回值避免结果堆积seen参数首次调用时初始化容量 50,000、误判率 0.0001 的 BloomFilter随后随子任务传递eventlet.Timeout(5, False)为请求设置 5 秒超时False表示超时后不抛异常而是返回None防止单个慢请求卡死整个绿色线程URL 提取使用宽松的正则来自 Daring Fireball 的 liberal URL 匹配并对每个匹配做同域校验location in domain(url)与布隆过滤器去重对每个新发现的 URL 构造子任务crawl.s(url, seen)聚合成 group 后delay()提交——注意这里是递归且是每层一个 group属于典型的扇出fan-out工作负载特别适合 Eventlet 这种以协程承载海量网络请求的模型。示例三bulk_task_producer.ProducerPool——连接复用的批量发布器要解决的问题日常客户端代码中如果每调用一次apply_async就新建一条 broker 连接在需要突发、大批量发布任务例如 Web 请求处理器内部时会产生巨大的连接开销。ProducerPool的解决方案是在单个进程内维护一个大小固定的 greenlet 池生产者每个 greenlet 只获取并复用一条自己的 producer/connection处理池内队列中的多个任务而不是每个任务都开新连接。使用方式模块文档字符串bulk_task_producer.py给出的用法 app Celery(brokeramqp://) pool ProducerPool(app, size20) receipt pool.apply_async(some_task, (1, 2), {}) receipt.wait() # block until the task has been published result receipt.result # the AsyncResult from task.apply_async调用方把(task, args, kwargs, options)元组推进内部队列池中任意空闲的 greenlet 取走并发布apply_async返回一个Receipt回执wait()阻塞直到任务已发布receipt.result拿到真正的AsyncResult。内部实现拆解Receipt发布回执——基于eventlet.event.Event实现的等待原语class Receipt: result None def __init__(self, callbackNone): self.callback callback self.ready Event() def finished(self, result): self.result result if self.callback: self.callback(result) self.ready.send() def wait(self, timeoutNone): with Timeout(timeout): return self.ready.wait()Event在 eventlet 中用于跨 greenlet 的信号通知任务真正发布成功后调用finished()设置结果并send()唤醒等待者wait()支持可选的超时参数。ProducerPool 主类class ProducerPool: Receipt Receipt def __init__(self, app, size20): self.app app self.size size self.inqueue LightQueue() self._running None self._producers None def apply_async(self, task, args, kwargs, callbackNone, **options): if self._running is None: self._running spawn_n(self._run) receipt self.Receipt(callback) self.inqueue.put((task, args, kwargs, options, receipt)) return receipt def _run(self): self._producers [ spawn_n(self._producer) for _ in range(self.size) ] def _producer(self): inqueue self.inqueue with self.app.producer_or_acquire() as producer: while 1: task, args, kwargs, options, receipt inqueue.get() result task.apply_async(args, kwargs, producerproducer, **options) receipt.finished(result)机制要点模块顶层调用monkey_patch()必须在任何网络/线程相关模块使用前完成补丁LightQueue是无界轻量队列greenlet 安全首次apply_async时spawn_n(self._run)懒启动随后_run派生size个_producergreenlet每个_producer通过self.app.producer_or_acquire()获取一个 producer——这正是 celery/app/base.py 中的连接池上下文管理器若未显式传入 producer则从producer_pool借出一个受broker_pool_limit、broker_pool_acquire_timeout等配置约束随后无限循环消费队列为每个任务复用同一条连接task.apply_async(..., producerproducer)因此整个批次只打开size默认 20条连接而不是每个任务一条——这正是 README 强调的以最小固定大小连接池换取突发批量发布吞吐的关键。源码纵深Celery 的 Eventlet TaskPool 是如何工作的示例运行时的执行池实现位于 celery/concurrency/eventlet.py。TaskPool(base.BasePool)的类属性标注了它的特性signal_safe False、is_green True、task_join_will_block False。底层是 eventlet GreenPoolon_start()中self._pool self.Pool(self.limit)其中limit即--concurrency指定的并发度示例中为 500所有任务通过GreenPool.spawn投递到池中执行每个任务一个 greenleton_apply()把目标任务包装为可被杀死的killable_target后 spawn并将id(greenlet)登记到_pool_map配合terminate_job()实现按任务终止kill能力补丁时机警告模块顶部有一段运行时检查若billiard.、celery.、kombu.等模块在 eventlet 补丁之前已经导入了thread/threading/socket会发出RuntimeWarningW_RACE提醒用户模块被过早加载、补丁失效——这与 celeryconfig.py 中不要用 worker_pool 配置项的告诫相互印证共同指向同一最佳实践尽量早地以命令行参数-P eventlet启动让补丁在加载 Celery 网络栈之前生效动态扩缩容grow(n)/shrink(n)直接调整GreenPool.size与信号量计数支持运行期调整并发度Eventlet 专用 TimerTimer类基于eventlet.greenthread.spawn_after实现延迟任务调度信号钩子on_start/on_stop会发送eventlet_pool_started、eventlet_pool_preshutdown、eventlet_pool_postshutdown等信号见celery.signals便于应用侧感知池生命周期。适用边界与注意事项结合 docs/userguide/concurrency/eventlet.rst 与示例代码使用 Eventlet 池时需要明确以下边界只适合 I/O 密集不要让单个任务长时间阻塞事件循环CPU 密集计算与 Eventlet 不契合绿色线程无法利用多核且计算会阻塞整个循环C 扩展库兼容性需逐一确认部分带 C 扩展的库无法被 monkey patch。文档举例pylibmc不能与 Eventlet 协作而psycopg2在特定条件下可以。使用第三方库前请查阅其文档patch 时机决定成败务必用-P eventlet命令行参数启动而非配置worker_pool否则补丁过晚、异步化失效参见 celery/concurrency/eventlet.py 的警告逻辑混合部署是常见方案官方文档建议可以同时运行 Eventlet 与 prefork 两类 worker并按任务类型兼容性或特性路由到不同的池结果后端限制示例中使用的iter_native()仅支持 amqp、Redis、cache 三类结果后端不要在生产中用示例爬虫直接全量外网抓取示例爬虫刻意限定同域名以防打爆外部站点去重过滤器按任务传递而非集中式高并发长跑场景应换用 Redis set 等集中方案。小结examples/eventlet目录用三个由浅入深的示例覆盖了 Celery Eventlet 的核心用法urlopen展示单任务与 group 并行调用的最小范式webcrawler.crawl展示递归扇出、布隆过滤器去重与 zlib 压缩序列化的工程组合bulk_task_producer.ProducerPool则示范了如何在客户端用固定大小 greenlet 池 连接复用突破突发批量发布的吞吐瓶颈。配合 celery/concurrency/eventlet.py 的 TaskPool 实现你可以在此基础上按需调整并发度、接入集中式去重服务或搭建 Eventlet prefork 混合 worker 集群。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价