ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

Python后端爬虫专题21:HTTP请求不能等爬虫跑完——Celery、Redis与Worker失败恢复 Python后端爬虫专题21HTTP请求不能等爬虫跑完——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,monthsparse_salary(str(raw[salary]))external_idstr(raw[external_id])returnJobItem(external_idexternal_id,source_urlf{base_url.rstrip(/)}/jobs/{external_id},titlestr(raw[title]),companystr(raw[company]),citystr(raw[city]),descriptionstr(raw[description]),skillslist(raw[skills]),salary_minsalary_min,salary_maxsalary_max,salary_monthsmonths,published_atstr(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],}jobparse_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-01HTTPX 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 转为可保存的 dicterrors 也变成普通字典。一次 Worker 崩溃的时间线task_acks_lateTrue表示 Worker 完成后才确认消息若进程中途死亡broker 可重新投递。task_reject_on_worker_lostTrue进一步要求丢失 Worker 时拒绝消息。重投会让同一任务执行两次所以第 10 篇的数据库唯一约束和幂等 upsert 是前提而不是优化项。这仍不是 exactly-once。Worker 可能已提交数据库但尚未 ack重跑时会得到 unchanged也可能在把任务状态设为 running 后崩溃数据库暂时显示 running。生产系统要有超时扫描识别长时间无心跳的任务并重试或标记消息队列不能自动修正业务表。哪些错误自动重试种子列表一页都没读到且存在失败时run_crawl_task抛RetryableCrawlTaskHTTP 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:int10)-CrawlReport:...classRetryableCrawlTask(RuntimeError):种子列表都无法读取队列可以在退避后重试整项任务。asyncdefrun_crawl_task(crawler:CrawlRunner,seed_url:str,*,tenant_id:str,max_pages:int10,)-dict[str,object]:执行一次采集并返回可由 JSON 序列化器保存的任务结果。reportawaitcrawler.run(seed_url,tenant_idtenant_id,max_pagesmax_pages)ifreport.list_pages0andreport.failed:reasonreport.errors[0].messageifreport.errorselseseed page failedraiseRetryableCrawlTask(reason)returnasdict(report)defcreate_celery_app(broker_url:str,result_backend:str|NoneNone)-Celery:建立安全的 Celery 序列化和 Worker 丢失恢复默认值。appCelery(jobradar,brokerbroker_url,backendresult_backendorbroker_url)app.conf.update(task_serializerjson,result_serializerjson,accept_content[json],task_acks_lateTrue,task_reject_on_worker_lostTrue,task_track_startedTrue,broker_connection_retry_on_startupTrue,)returnapp本篇课后练习写出消息在“收到前、执行中、数据库已提交但未 ack”三个时点 Worker 崩溃后的结果说明为什么会重复以及幂等如何收敛。解释 Redis broker、Celery result backend、PostgreSQL crawl_tasks 各保存什么为什么不能互相替代。将CrawlReport(list_pages1, created2)交给run_crawl_task写出完整异步脚本并证明返回值可json.dumps。下一篇将从 FastAPI 创建任务并查询结果。
返回列表