
Python后端爬虫专题22从一个202响应开始——FastAPI任务创建、状态查询与租户隔离上一篇练习完整答案Worker 收到消息前崩溃Redis 中未确认消息仍等待其他 Worker执行中崩溃由于 late ack 和 reject_on_worker_lost 会重投数据库已提交但尚未 ack 时崩溃也会重投第二次运行依靠(tenant_id, source_url)唯一约束和内容指纹变为 unchanged而非新增职位。Exactly-once 不能靠一句配置获得。Redis broker 保存待执行消息与投递状态result backend 保存 Celery 层结果通常有过期策略PostgreSQL crawl_tasks 保存面向用户的业务状态、租户、seed、报告和错误。前两者服务任务基础设施后一项服务产品合同和审计。完整 JSON 序列化验证importasyncioimportjsonfromjobradar.pipelineimportCrawlReportfromjobradar.tasksimportrun_crawl_taskclassDemoCrawler:asyncdefrun(self,seed_url:str,*,tenant_id:str,max_pages:int10):returnCrawlReport(list_pages1,discovered2,created2)resultasyncio.run(run_crawl_task(DemoCrawler(),http://target.test/jobs,tenant_idtenant-a))encodedjson.dumps(result,ensure_asciiFalse)assertcreated: 2inencoded先从调用者视角看四个接口POST /api/crawls接收 seed_url 与 max_pages要求X-Tenant-ID返回 202、业务 id、queued 与 queue_task_idGET /api/crawls/{id}返回状态和报告GET /api/jobs?page1size20返回职位分页GET /api/analytics返回聚合。健康检查不要求租户只说明 API 进程活着不代表 Redis、Worker 和采集链都正常。用 curl 创建课程任务$headers {X-Tenant-IDtenant-a}$body {seed_urlhttp://targetlab:8011/jobs?page1;max_pages2}|ConvertTo-Json$taskInvoke-RestMethod-Method Post-Uri http://127.0.0.1:8010/api/crawls -Headers$headers-ContentType application/json-Body$body$taskInvoke-RestMethod-Urihttp://127.0.0.1:8010/api/crawls/$($task.id)-Headers$headerstargetlab是 Compose 网络内名称在宿主机直接访问页面用 127.0.0.1:8011但真正抓取者是容器 Worker所以种子必须是它可解析的地址。请求经过哪些验证Pydantic 先检查绝对 HTTP(S) 形式和 max_pages 1—100依赖函数清理 X-Tenant-ID 并拒绝空白CrawlPolicy 再检查 host allowlist 与私网授权。验证必须在写库和入队之前否则恶意 seed 即使最终失败也已经把服务器变成 URL 探测器。创建流程是生成业务 UUIDRepository 建 queued 记录并 commit队列入队回写 queue_task_id 并 commit。入队异常写 dispatch_failedAPI 返回 503。为什么中间多次 commit让已接受任务在队列故障时仍留下可诊断事实。代价是存在第 11 篇讨论的 outbox 窗口。cd project.\.venv\Scripts\python.exe-m pytest tests\test_api.py::test_create_crawl_validates_seed_persists_task_and_enqueues_once-q测试不只看 202。它断言队列恰好收到一次四字段消息再用真实 Repository 查询任务和 queue_task_id。若路由只返回一个漂亮 JSON、实际没有落库或投递检查会失败。为什么跨租户查询返回404任务 ID 即使是 UUID也不应视为权限。Repository 查询同时带 tenant_idtenant-b 请求 tenant-a 任务时返回 404而非 403避免泄露该 ID 确实存在。职位列表和统计也必须在 SQL 查询阶段过滤租户不能先查全部再在 Python 中剔除。课程用请求头模拟上游认证已经确定的租户。生产中不能让公网用户任意填写 X-Tenant-ID而应由 JWT/API key 映射、网关签名头或内部身份系统注入。这里教学重点是每层都显式传递 tenant而不是实现完整账号系统。分页要约束查询不只改响应元数据offset(page-1)*size、limitsize传给 Repositorytotal 独立 count。若只在响应中写 page2却仍返回全部行前端很快发现重复。本项目按 id 稳定排序职位很多时 offset 深分页变慢可改基于(id cursor)的游标分页。size 上限 100 防止一个请求把所有描述载入内存。统计接口目前为教学简化一次读取 100000 行实际大数据应让数据库做 GROUP BY 或建立离线聚合下一篇会明确这一边界。OpenAPI能替代教程吗FastAPI 自动提供/docs可看到字段和状态码却看不到“先落库再入队”的故障窗口、租户信任来源和为什么 304 不写空快照。OpenAPI 描述接口形状文章解释业务语义两者都要有。本篇完整 FastAPI 模块阅读顺序请求模型、对外职位 payload、create_app 的依赖装配、创建任务、查任务、分页、统计。路由只做边界和编排采集逻辑不在请求线程执行。FastAPI 边界创建采集任务并按租户查询任务与职位。fromcollections.abcimportCallablefromtypingimportAnnotated,Protocolfromuuidimportuuid4fromfastapiimportDepends,FastAPI,Header,HTTPException,Query,statusfrompydanticimportBaseModel,Fieldfromsqlalchemy.ormimportSessionfrom.analyticsimportsummarize_jobsfrom.policyimportCrawlPolicy,PolicyViolationfrom.repositoryimportJobRecord,JobRepositoryclassTaskQueue(Protocol):defenqueue(self,*,task_id:str,seed_url:str,tenant_id:str,max_pages:int)-str:...classCreateCrawlRequest(BaseModel):seed_url:strField(patternr^https?://,max_length2048)max_pages:intField(default3,ge1,le100)def_job_payload(row:JobRecord)-dict[str,object]:return{id:row.id,external_id:row.external_id,source_url:row.source_url,title:row.title,company:row.company,city:row.city,description:row.description,skills:row.skills,salary_min:row.salary_min,salary_max:row.salary_max,salary_months:row.salary_months,published_at:row.published_at.isoformat(),}defcreate_app(session_factory:Callable[[],Session],queue:TaskQueue,*,allowed_seed_hosts:set[str],allowed_private_hosts:set[str]|NoneNone,)-FastAPI:装配 API调用者显式提供数据库和队列测试不会连接生产服务。appFastAPI(titleJobRadar API,version0.1.0)policyCrawlPolicy(allowed_hostsallowed_seed_hosts,allowed_private_hostsallowed_private_hostsorset(),)deftenant_id(x_tenant_id:Annotated[str,Header(aliasX-Tenant-ID,min_length1)],)-str:cleanedx_tenant_id.strip()ifnotcleaned:raiseHTTPException(status_code422,detailX-Tenant-ID must not be blank)returncleanedapp.get(/health)defhealth()-dict[str,str]:return{status:ok,service:jobradar-api}app.post(/api/crawls,status_codestatus.HTTP_202_ACCEPTED)defcreate_crawl(request:CreateCrawlRequest,current_tenant:Annotated[str,Depends(tenant_id)],)-dict[str,object]:try:policy.check_url(request.seed_url)exceptPolicyViolationasexc:raiseHTTPException(status_code422,detailstr(exc))fromexc task_idstr(uuid4())withsession_factory()assession:repositoryJobRepository(session)repository.create_crawl_task(task_id,tenant_idcurrent_tenant,seed_urlrequest.seed_url,max_pagesrequest.max_pages,)repository.commit()try:queue_task_idqueue.enqueue(task_idtask_id,seed_urlrequest.seed_url,tenant_idcurrent_tenant,max_pagesrequest.max_pages,)exceptExceptionasexc:repository.mark_crawl_task_failed(current_tenant,task_id,type(exc).__name__)repository.commit()raiseHTTPException(status_code503,detailtask queue unavailable)fromexc repository.attach_queue_task(current_tenant,task_id,queue_task_id)repository.commit()return{id:task_id,status:queued,queue_task_id:queue_task_id,}app.get(/api/crawls/{task_id})defget_crawl(task_id:str,current_tenant:Annotated[str,Depends(tenant_id)],)-dict[str,object]:withsession_factory()assession:taskJobRepository(session).get_crawl_task(current_tenant,task_id)iftaskisNone:raiseHTTPException(status_code404,detailcrawl task not found)return{id:task.id,seed_url:task.seed_url,max_pages:task.max_pages,status:task.status,queue_task_id:task.queue_task_id,error_message:task.error_message,report:task.report,}app.get(/api/jobs)deflist_jobs(current_tenant:Annotated[str,Depends(tenant_id)],page:Annotated[int,Query(ge1)]1,size:Annotated[int,Query(ge1,le100)]20,)-dict[str,object]:withsession_factory()assession:repositoryJobRepository(session)rowsrepository.list_jobs(current_tenant,offset(page-1)*size,limitsize)totalrepository.count_jobs(current_tenant)return{items:[_job_payload(row)forrowinrows],page:page,size:size,total:total,}app.get(/api/analytics)defanalytics(current_tenant:Annotated[str,Depends(tenant_id)],top_n:Annotated[int,Query(ge1,le50)]10,)-dict[str,object]:withsession_factory()assession:rowsJobRepository(session).list_jobs(current_tenant,limit100_000)summarysummarize_jobs(rows,top_ntop_n)return{total_jobs:summary.total_jobs,city_counts:summary.city_counts,top_skills:[{skill:skill,count:count}forskill,countinsummary.top_skills],salary_sample_size:summary.salary_sample_size,average_annual_salary:summary.average_annual_salary,}returnapp本篇课后练习写出创建任务成功、seed 不在 allowlist、队列不可用、跨租户查询四个场景的 HTTP 状态与数据库/队列副作用。在 ASGITransport 中创建三条职位请求page2size2给出完整测试并断言只有第三条。解释为什么客户端提供的 X-Tenant-ID 在生产中不能直接信任给出两种可信注入方法。下一篇会把职位变成可解释的技能、城市与薪资趋势。