news 2026/10/3 2:26:56

Python后端爬虫专题21:HTTP请求不能等爬虫跑完——Celery、Redis与Worker失败恢复

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Python后端爬虫专题21:HTTP请求不能等爬虫跑完——Celery、Redis与Worker失败恢复

Python后端爬虫专题21:HTTP请求不能等爬虫跑完——Celery、Redis与Worker失败恢复

上一篇练习完整答案

四类场景的选择是:初始 HTML 已含完整业务节点,选 HTTPX + HTML parser,遇到结构变化明确失败;HTML 无数据但有稳定、获准且字段完整的 JSON,选 HTTPX + JSON adapter,接口返回 401/403/合同外字段时停止;必须执行脚本才得到获准数据,选 Playwright + shared parser,遇验证码或超出授权交互时停止;没有授权则不实现采集,联系负责人获得 API/账号/频率范围。

TargetLab JSON 的完整适配示例:

fromjobradar.modelsimportJobItem,parse_salarydefparse_api_job(raw:dict[str,object],base_url:str)->JobItem:salary_min,salary_max,months=parse_salary(str(raw["salary"]))external_id=str(raw["external_id"])returnJobItem(external_id=external_id,source_url=f"{base_url.rstrip('/')}/jobs/{external_id}",title=str(raw["title"]),company=str(raw["company"]),city=str(raw["city"]),description=str(raw["description"]),skills=list(raw["skills"]),salary_min=salary_min,salary_max=salary_max,salary_months=months,published_at=str(raw["published_at"]),)first={"external_id":"python-backend-001","title":"Python 后端工程师","company":"星河科技","city":"杭州","salary":"15k-25k·13薪","published_at":"2026-09-01","description":"负责后端平台。","skills":["Python","FastAPI","PostgreSQL"],}job=parse_api_job(first,"http://target.test")assert(job.salary_min,job.salary_max,job.salary_months)==(15000,25000,13)assertjob.published_at.isoformat()=="2026-09-01"

HTTPX API 路线资源少、结构化、容易契约测试,但可能不是正式接口或与页面字段不同;Playwright CPU/内存和运维成本高、波动大,却能得到获准的最终渲染结果。选型顺序是授权与稳定合同、字段正确性、可测试性,最后才是性能。

为什么 POST /crawls 应该返回 202

一次采集可能包含分页、数百详情、重试和数据库写入。如果 FastAPI 一直等完成,客户端连接超时、Web Worker 被长期占用,重试 HTTP 请求还可能重复创建任务。JobRadar 接收请求后验证 seed、创建业务任务、投递 Celery,然后以202 Accepted返回 task id;客户端轮询任务接口。

Redis 在这里是 broker:保存待消费消息并把它交给 Worker。它不是职位事实数据库;PostgreSQL 才保存 crawl_tasks 和 jobs。Celery result backend 也不能替代业务任务表,因为结果会过期、任务迁移时 ID 可能改变、租户访问控制也属于应用。

消息里只传四个标量

API 入队参数是task_id、seed_url、tenant_id、max_pages。不能传 SQLAlchemy Session、HttpFetcher、Crawler 或 Pydantic 对象:它们包含连接与进程状态,无法安全跨进程序列化。Worker 收到标量后从环境构造数据库、快照存储、策略和 HTTP client。

Celery 配置只接受 JSON,明确拒绝 pickle。pickle 可以在反序列化时执行任意代码,不应让不可信消息触发。run_crawl_task最终用asdict把 CrawlReport 转为可保存的 dict,errors 也变成普通字典。

一次 Worker 崩溃的时间线

task_acks_late=True表示 Worker 完成后才确认消息;若进程中途死亡,broker 可重新投递。task_reject_on_worker_lost=True进一步要求丢失 Worker 时拒绝消息。重投会让同一任务执行两次,所以第 10 篇的数据库唯一约束和幂等 upsert 是前提,而不是优化项。

这仍不是 exactly-once。Worker 可能已提交数据库但尚未 ack,重跑时会得到 unchanged;也可能在把任务状态设为 running 后崩溃,数据库暂时显示 running。生产系统要有超时扫描,识别长时间无心跳的任务并重试或标记;消息队列不能自动修正业务表。

哪些错误自动重试

种子列表一页都没读到且存在失败时,run_crawl_task抛RetryableCrawlTask;HTTP transport 断连也可重试。Celery 最多重试三次,使用 backoff 与 jitter,避免许多 Worker 同时打回故障服务。详情里一条解析失败则保留 partial 报告,不重跑整批,否则每个永久坏页面会反复拖累任务。

cd project.\.venv\Scripts\python.exe-m pytest tests\test_tasks.py::test_run_crawl_task_marks_seed_failure_as_retryable-q.\.venv\Scripts\python.exe-m pytest tests\test_tasks.py::test_celery_app_uses_json_and_late_acknowledgement-q

第一个检查业务失败分类,第二个检查消息安全与确认策略。只测试 Celery decorator 存在没有价值,必须验证会影响恢复行为的配置。

业务ID与Celery ID为何分开

API 先生成 UUID 作为业务 task_id,写入 PostgreSQL,再调用 Celery 得到 queue_task_id。业务 ID 面向 API、租户和审计,重投也可保持;队列 ID 面向 broker/result backend。把两者混用后,迁移队列或重新投递会让客户端旧链接失效。

先落库再投递会留下“已排队但未发送”的崩溃窗口,第 11 篇已经用 outbox 解释。当前实现投递异常会写dispatch_failed并返回 503,至少错误可见;不是把两次操作伪装成原子事务。

本篇完整任务边界模块

tasks.py不创建真实数据库或 HTTP client,只定义可测试的应用边界与 Celery 默认值。具体生产装配在worker.py,这样单测无需启动 Redis。

"""采集应用服务与 Celery 之间的任务边界。"""fromdataclassesimportasdictfromtypingimportProtocolfromceleryimportCeleryfrom.pipelineimportCrawlReportclassCrawlRunner(Protocol):asyncdefrun(self,seed_url:str,*,tenant_id:str,max_pages:int=10)->CrawlReport:...classRetryableCrawlTask(RuntimeError):"""种子列表都无法读取;队列可以在退避后重试整项任务。"""asyncdefrun_crawl_task(crawler:CrawlRunner,seed_url:str,*,tenant_id:str,max_pages:int=10,)->dict[str,object]:"""执行一次采集并返回可由 JSON 序列化器保存的任务结果。"""report=awaitcrawler.run(seed_url,tenant_id=tenant_id,max_pages=max_pages)ifreport.list_pages==0andreport.failed:reason=report.errors[0].messageifreport.errorselse"seed page failed"raiseRetryableCrawlTask(reason)returnasdict(report)defcreate_celery_app(broker_url:str,result_backend:str|None=None)->Celery:"""建立安全的 Celery 序列化和 Worker 丢失恢复默认值。"""app=Celery("jobradar",broker=broker_url,backend=result_backendorbroker_url)app.conf.update(task_serializer="json",result_serializer="json",accept_content=["json"],task_acks_late=True,task_reject_on_worker_lost=True,task_track_started=True,broker_connection_retry_on_startup=True,)returnapp

本篇课后练习

  1. 写出消息在“收到前、执行中、数据库已提交但未 ack”三个时点 Worker 崩溃后的结果,说明为什么会重复以及幂等如何收敛。
  2. 解释 Redis broker、Celery result backend、PostgreSQL crawl_tasks 各保存什么,为什么不能互相替代。
  3. 将CrawlReport(list_pages=1, created=2)交给run_crawl_task,写出完整异步脚本并证明返回值可json.dumps。下一篇将从 FastAPI 创建任务并查询结果。
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/3 2:26:41

从零实现 mini-git:用真实 Git 验证 blob、tree、commit 和 index

从零实现 mini-git:用真实 Git 验证 blob、tree、commit 和 index项目地址:https://github.com/yituanxing/mini-git我一开始写 mini-git,不是为了再造一个能替代 Git 的工具。真正的动机更简单:很多 Git 概念背起来都像八股&…

作者头像 李华
网站建设 2026/10/3 2:24:37

darktable:免费实用的RAW处理器,快速从导入到出片

darktable:免费实用的RAW处理器,快速从导入到出片 【免费下载链接】darktable darktable is an open source photography workflow application and raw developer 项目地址: https://gitcode.com/GitHub_Trending/da/darktable RAW 打开灰蒙蒙&a…

作者头像 李华
网站建设 2026/10/3 2:24:36

一次扫 100 个仓库会不会失控?Codex Security 批量扫描成本门禁

一次扫 100 个仓库会不会失控?Codex Security 批量扫描成本门禁 [!NOTE] Codex Security 的 bulk-scan 能发现 GitHub 仓库或读取固定提交的 CSV,并把每个仓库隔离成可续跑、可审计的独立尝试。 真正的成本门禁不是把 workers 调大,而是先冻结清单、分波次试点、区分仓库级并…

作者头像 李华