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本篇课后练习
- 写出消息在“收到前、执行中、数据库已提交但未 ack”三个时点 Worker 崩溃后的结果,说明为什么会重复以及幂等如何收敛。
- 解释 Redis broker、Celery result backend、PostgreSQL crawl_tasks 各保存什么,为什么不能互相替代。
- 将
CrawlReport(list_pages=1, created=2)交给run_crawl_task,写出完整异步脚本并证明返回值可json.dumps。下一篇将从 FastAPI 创建任务并查询结果。