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:int=10):returnCrawlReport(list_pages=1,discovered=2,created=2)result=asyncio.run(run_crawl_task(DemoCrawler(),"http://target.test/jobs",tenant_id="tenant-a"))encoded=json.dumps(result,ensure_ascii=False)assert'"created": 2'inencoded先从调用者视角看四个接口
POST /api/crawls接收 seed_url 与 max_pages,要求X-Tenant-ID,返回 202、业务 id、queued 与 queue_task_id;GET /api/crawls/{id}返回状态和报告;GET /api/jobs?page=1&size=20返回职位分页;GET /api/analytics返回聚合。健康检查不要求租户,只说明 API 进程活着,不代表 Redis、Worker 和采集链都正常。
用 curl 创建课程任务:
$headers= @{"X-Tenant-ID"="tenant-a"}$body= @{seed_url="http://targetlab:8011/jobs?page=1";max_pages=2}|ConvertTo-Json$task=Invoke-RestMethod-Method Post-Uri http://127.0.0.1:8010/api/crawls `-Headers$headers-ContentType application/json-Body$body$taskInvoke-RestMethod-Uri"http://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 探测器。
创建流程是:生成业务 UUID;Repository 建 queued 记录并 commit;队列入队;回写 queue_task_id 并 commit。入队异常写 dispatch_failed,API 返回 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_id;tenant-b 请求 tenant-a 任务时返回 404,而非 403,避免泄露该 ID 确实存在。职位列表和统计也必须在 SQL 查询阶段过滤租户,不能先查全部再在 Python 中剔除。
课程用请求头模拟上游认证已经确定的租户。生产中不能让公网用户任意填写 X-Tenant-ID,而应由 JWT/API key 映射、网关签名头或内部身份系统注入。这里教学重点是每层都显式传递 tenant,而不是实现完整账号系统。
分页要约束查询,不只改响应元数据
offset=(page-1)*size、limit=size传给 Repository,total 独立 count。若只在响应中写 page=2,却仍返回全部行,前端很快发现重复。本项目按 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:str=Field(pattern=r"^https?://",max_length=2048)max_pages:int=Field(default=3,ge=1,le=100)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]|None=None,)->FastAPI:"""装配 API;调用者显式提供数据库和队列,测试不会连接生产服务。"""app=FastAPI(title="JobRadar API",version="0.1.0")policy=CrawlPolicy(allowed_hosts=allowed_seed_hosts,allowed_private_hosts=allowed_private_hostsorset(),)deftenant_id(x_tenant_id:Annotated[str,Header(alias="X-Tenant-ID",min_length=1)],)->str:cleaned=x_tenant_id.strip()ifnotcleaned:raiseHTTPException(status_code=422,detail="X-Tenant-ID must not be blank")returncleaned@app.get("/health")defhealth()->dict[str,str]:return{"status":"ok","service":"jobradar-api"}@app.post("/api/crawls",status_code=status.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_code=422,detail=str(exc))fromexc task_id=str(uuid4())withsession_factory()assession:repository=JobRepository(session)repository.create_crawl_task(task_id,tenant_id=current_tenant,seed_url=request.seed_url,max_pages=request.max_pages,)repository.commit()try:queue_task_id=queue.enqueue(task_id=task_id,seed_url=request.seed_url,tenant_id=current_tenant,max_pages=request.max_pages,)exceptExceptionasexc:repository.mark_crawl_task_failed(current_tenant,task_id,type(exc).__name__)repository.commit()raiseHTTPException(status_code=503,detail="task 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:task=JobRepository(session).get_crawl_task(current_tenant,task_id)iftaskisNone:raiseHTTPException(status_code=404,detail="crawl 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(ge=1)]=1,size:Annotated[int,Query(ge=1,le=100)]=20,)->dict[str,object]:withsession_factory()assession:repository=JobRepository(session)rows=repository.list_jobs(current_tenant,offset=(page-1)*size,limit=size)total=repository.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(ge=1,le=50)]=10,)->dict[str,object]:withsession_factory()assession:rows=JobRepository(session).list_jobs(current_tenant,limit=100_000)summary=summarize_jobs(rows,top_n=top_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 中创建三条职位,请求
page=2&size=2,给出完整测试并断言只有第三条。 - 解释为什么客户端提供的 X-Tenant-ID 在生产中不能直接信任,给出两种可信注入方法。下一篇会把职位变成可解释的技能、城市与薪资趋势。