Python后端爬虫专题11:别只建一张jobs表——SQLAlchemy模型、任务状态与事务边界
上一篇练习完整答案
指纹实验应得到:技能换序/大小写不变,正文变化后不同。last_seen_at每次任务都会更新,若进入指纹,unchanged 永远不会出现。并发插入时唯一约束可能让后提交者收到 IntegrityError;正确处理是回滚该事务后重读现有行,而不是删约束。
fromjobradar.fingerprintimportjob_fingerprint first=job_item same=first.model_copy(update={"skills":["PostgreSQL","python","FastAPI"]})changed=first.model_copy(update={"description":"新的职责"})assertjob_fingerprint(first)==job_fingerprint(same)assertjob_fingerprint(first)!=job_fingerprint(changed)任务与职位是两类事实
jobs 回答“当前有哪些职位”;crawl_tasks 回答“谁在什么时候请求了什么采集、当前进行到哪里”。只建 jobs 表,API 返回 202 后用户无法查询任务是排队、运行、部分成功还是失败;只依赖 Celery result backend,业务任务又容易随过期策略消失。
CrawlTaskRecord 保存 UUID、tenant、seed、max_pages、status、queue_task_id、error_message 和 report。queue_task_id 是中间件身份,不替代业务 task id:重试或迁移队列时业务任务仍能保持稳定。
状态不是随便写字符串
课程允许 queued、running、completed、partial、failed、dispatch_failed。创建时先 queued 并提交,再调用队列;入队失败就标记 dispatch_failed。Worker 开始执行时改 running,报告有局部失败则 partial,种子无法采集或进程异常为 failed。
更严格的生产实现还应约束状态转换,例如 completed 不能回到 running,并用版本号避免 API 与 Worker 并发覆盖。课程先通过 repository 集中入口,而不是在路由中直接改属性,为以后添加转换规则留出位置。
为什么先落库再入队
如果先入队,Worker 可能在 API 提交任务记录前开始,查询不到 task_id;如果先提交数据库再入队,队列失败会留下 dispatch_failed,可由后台扫描重试。这仍不是完全原子:进程可能在提交后、入队前崩溃。可靠生产系统可采用 transactional outbox,由同一数据库事务写任务和待发布事件。
课程没有为了显得“企业级”直接塞入 outbox 全套,因为本专题核心是采集;但文章必须诚实说明故障窗口,而不是把两次操作说成事务。
Session 生命周期与 commit
API 每次请求with session_factory()创建 Session;Crawler 每个详情 upsert 后 commit,使坏详情不回滚此前成功记录;异常路径 rollback 后才能继续使用 Session。Session 不是线程安全对象,不能传进 Celery 消息或跨协程任意共享。
数据库 datetime 使用 UTC aware 值,显示时再转换时区。JSON report 只存可序列化标量,不能放 PageParseError、Pydantic 模型或 SQLAlchemy 对象。
Alembic 与 create_all 的分工
本地测试create_schema()快速从 metadata 建空库;Compose 设置 auto_create_schema=false,由alembic upgrade head管理迁移版本。生产不能每次启动盲目 create_all,因为它不会可靠修改已有列,也没有可审计升级路径。
迁移测试在全新 SQLite 文件执行 Alembic,检查 alembic_version、crawl_tasks、jobs 三张表。这保证迁移脚本不是摆设。PostgreSQL 在第 27 篇用 Compose 再验证。
本篇行为检查
.\.venv\Scripts\python.exe-m pytest tests\test_repository.py::test_repository_records_worker_status_and_report_without_cross_tenant_update-q测试创建任务、running、completed 并保存报告,然后用 tenant-b 查询得到 None。若 repository 只按 task id 更新,跨租户风险会被抓到。
本篇完整 Repository
文件较长,请按四块阅读:ORM 表定义、UpsertResult、职位方法、任务方法。_item_columns只做 JobItem 到列的集中映射,真正由 upsert 调用;它不是为了测试存在的展示 helper。
"""职位持久化:唯一约束、内容变化和租户隔离集中在这里。"""fromdataclassesimportdataclassfromdatetimeimportdate,datetime,timezonefromsqlalchemyimportDate,DateTime,Integer,JSON,String,Text,UniqueConstraint,func,selectfromsqlalchemy.engineimportEnginefromsqlalchemy.ormimportDeclarativeBase,Mapped,Session,mapped_columnfrom.fingerprintimportjob_fingerprintfrom.modelsimportJobItemdefutc_now()->datetime:returndatetime.now(timezone.utc)classBase(DeclarativeBase):passclassJobRecord(Base):__tablename__="jobs"__table_args__=(UniqueConstraint("tenant_id","source_url",name="uq_jobs_tenant_source"),)id:Mapped[int]=mapped_column(primary_key=True,autoincrement=True)tenant_id:Mapped[str]=mapped_column(String(80),index=True)external_id:Mapped[str]=mapped_column(String(120))source_url:Mapped[str]=mapped_column(String(2048))title:Mapped[str]=mapped_column(String(200))company:Mapped[str]=mapped_column(String(200))city:Mapped[str]=mapped_column(String(100),index=True)description:Mapped[str]=mapped_column(Text)skills:Mapped[list[str]]=mapped_column(JSON,default=list)salary_min:Mapped[int|None]=mapped_column(Integer,nullable=True)salary_max:Mapped[int|None]=mapped_column(Integer,nullable=True)salary_months:Mapped[int]=mapped_column(Integer,default=12)published_at:Mapped[date]=mapped_column(Date)content_fingerprint:Mapped[str]=mapped_column(String(64))etag:Mapped[str|None]=mapped_column(String(255),nullable=True)last_modified:Mapped[str|None]=mapped_column(String(255),nullable=True)snapshot_id:Mapped[str|None]=mapped_column(String(64),nullable=True)first_seen_at:Mapped[datetime]=mapped_column(DateTime(timezone=True),default=utc_now)last_seen_at:Mapped[datetime]=mapped_column(DateTime(timezone=True),default=utc_now)updated_at:Mapped[datetime]=mapped_column(DateTime(timezone=True),default=utc_now)classCrawlTaskRecord(Base):__tablename__="crawl_tasks"id:Mapped[str]=mapped_column(String(36),primary_key=True)tenant_id:Mapped[str]=mapped_column(String(80),index=True)seed_url:Mapped[str]=mapped_column(String(2048))max_pages:Mapped[int]=mapped_column(Integer)status:Mapped[str]=mapped_column(String(30),default="queued")queue_task_id:Mapped[str|None]=mapped_column(String(255),nullable=True)error_message:Mapped[str|None]=mapped_column(Text,nullable=True)report:Mapped[dict[str,object]|None]=mapped_column(JSON,nullable=True)created_at:Mapped[datetime]=mapped_column(DateTime(timezone=True),default=utc_now)updated_at:Mapped[datetime]=mapped_column(DateTime(timezone=True),default=utc_now)@dataclass(frozen=True)classUpsertResult:job_id:intaction:strdefcreate_schema(engine:Engine)->None:"""只用于本地练习和测试;生产环境由 Alembic 迁移建表。"""Base.metadata.create_all(engine)classJobRepository:"""封装 JobRecord 查询,让流水线不拼 SQLAlchemy 语句。"""def__init__(self,session:Session)->None:self._session=sessiondefupsert(self,tenant_id:str,item:JobItem,*,etag:str|None=None,last_modified:str|None=None,snapshot_id:str|None=None,)->UpsertResult:fingerprint=job_fingerprint(item)record=self._session.scalar(select(JobRecord).where(JobRecord.tenant_id==tenant_id,JobRecord.source_url==item.source_url,))now=utc_now()ifrecordisNone:record=JobRecord(tenant_id=tenant_id,content_fingerprint=fingerprint,first_seen_at=now,last_seen_at=now,updated_at=now,**_item_columns(item),)record.etag=etag record.last_modified=last_modified record.snapshot_id=snapshot_id self._session.add(record)self._session.flush()returnUpsertResult(job_id=record.id,action="created")action="unchanged"ifrecord.content_fingerprint!=fingerprint:forname,valuein_item_columns(item).items():setattr(record,name,value)record.content_fingerprint=fingerprint record.updated_at=now action="updated"record.last_seen_at=now record.etag=etagorrecord.etag record.last_modified=last_modifiedorrecord.last_modified record.snapshot_id=snapshot_idorrecord.snapshot_id self._session.flush()returnUpsertResult(job_id=record.id,action=action)defget_validators(self,tenant_id:str,source_url:str)->tuple[str|None,str|None]:row=self._session.execute(select(JobRecord.etag,JobRecord.last_modified).where(JobRecord.tenant_id==tenant_id,JobRecord.source_url==source_url,)).one_or_none()returnrowifrowisnotNoneelse(None,None)deflist_jobs(self,tenant_id:str,*,offset:int=0,limit:int=100)->list[JobRecord]:statement=(select(JobRecord).where(JobRecord.tenant_id==tenant_id).order_by(JobRecord.id).offset(offset).limit(limit))returnlist(self._session.scalars(statement))defcount_jobs(self,tenant_id:str)->int:returnint(self._session.scalar(select(func.count()).select_from(JobRecord).where(JobRecord.tenant_id==tenant_id))or0)defcreate_crawl_task(self,task_id:str,*,tenant_id:str,seed_url:str,max_pages:int,)->CrawlTaskRecord:task=CrawlTaskRecord(id=task_id,tenant_id=tenant_id,seed_url=seed_url,max_pages=max_pages,status="queued",)self._session.add(task)self._session.flush()returntaskdefattach_queue_task(self,tenant_id:str,task_id:str,queue_task_id:str)->None:task=self.get_crawl_task(tenant_id,task_id)iftaskisNone:raiseLookupError("crawl task not found")task.queue_task_id=queue_task_id task.updated_at=utc_now()self._session.flush()defmark_crawl_task_failed(self,tenant_id:str,task_id:str,message:str)->None:task=self.get_crawl_task(tenant_id,task_id)iftaskisNone:raiseLookupError("crawl task not found")task.status="dispatch_failed"task.error_message=message[:1000]task.updated_at=utc_now()self._session.flush()defset_crawl_task_status(self,tenant_id:str,task_id:str,status:str,*,report:dict[str,object]|None=None,error_message:str|None=None,)->None:allowed={"queued","running","completed","partial","failed","dispatch_failed"}ifstatusnotinallowed:raiseValueError(f"unsupported crawl task status:{status}")task=self.get_crawl_task(tenant_id,task_id)iftaskisNone:raiseLookupError("crawl task not found")task.status=status task.report=report task.error_message=error_message[:1000]iferror_messageelseNonetask.updated_at=utc_now()self._session.flush()defget_crawl_task(self,tenant_id:str,task_id:str)->CrawlTaskRecord|None:returnself._session.scalar(select(CrawlTaskRecord).where(CrawlTaskRecord.tenant_id==tenant_id,CrawlTaskRecord.id==task_id,))defcommit(self)->None:self._session.commit()defrollback(self)->None:self._session.rollback()def_item_columns(item:JobItem)->dict[str,object]:return{"external_id":item.external_id,"source_url":item.source_url,"title":item.title,"company":item.company,"city":item.city,"description":item.description,"skills":list(item.skills),"salary_min":item.salary_min,"salary_max":item.salary_max,"salary_months":item.salary_months,"published_at":item.published_at,}本篇课后练习
- 画出 queued→running→completed/partial/failed 状态图,并标出 API 与 Worker 各负责哪些边。
- 解释“先落库再入队”仍有哪些崩溃窗口,提出 outbox 的最小表结构。
- 在临时库运行 Alembic,再用 SQLAlchemy inspect 列出三张表。下一篇会让每个解析结果都能追溯到原始 HTML。