6.1 Item Pipeline:从清洗到落库 本节摘要:Item Pipeline 是回程行李的分拣流水线:每个 Item 依次穿过启用的管道,清洗、校验、去重、落库各占一层。本节讲管道的生命周期钩子、执行顺序规则与分层设计模式,并把数据清洗的标准动作(去空白、类型归一、缺失处理)一次配齐。 Spider 回调 yield 出 Item 的那一刻,回程开始。第 2 章给数据上了户口(契约模型),本节按契约对货物验收加工。管道是数据质量的最后防线——它之外的任何环节出错都还有救,脏数据一旦落库就只剩返工。 生命周期与顺序规则 管道是普通类加三个钩子:processitem 干活,openspider 与 closespider 管开收工。
本节摘要:Item Pipeline 是回程行李的分拣流水线:每个 Item 依次穿过启用的管道,清洗、校验、去重、落库各占一层。本节讲管道的生命周期钩子、执行顺序规则与分层设计模式,并把数据清洗的标准动作(去空白、类型归一、缺失处理)一次配齐。
Spider 回调 yield 出 Item 的那一刻,回程开始。第 2 章给数据上了户口(契约模型),本节按契约对货物验收加工。管道是数据质量的最后防线——它之外的任何环节出错都还有救,脏数据一旦落库就只剩返工。
管道是普通类加三个钩子:process_item 干活,open_spider 与 close_spider 管开收工。执行顺序由 settings 里装配数字决定,数字小的先跑:
class CleanPipeline: def open_spider(self, spider): # 第一道管道开工时调用:适合准备共享资源 self.seen_keys = set() def process_item(self, item, spider): # 每个 Item 依次路过:返回 item 继续走,抛异常则中断本件 return item def close_spider(self, spider): # 收工:落日志、关连接 spider.logger.info("清洗管道共经手 %s 件", len(self.seen_keys))
# settings 配置模块:管道分层的顺序即设计 ITEM_PIPELINES = { "bookstation.pipelines.CleanPipeline": 100, # 先洗 "bookstation.pipelines.ValidatePipeline": 200, # 再验 "bookstation.pipelines.DedupPipeline": 300, # 后去重 "bookstation.pipelines.MySqlPipeline": 400, # 最后落库 }
顺序设计的原则只有一条:越便宜的检查越靠前。清洗最便宜放最前,落库最贵放最后——让废件尽早被拦下,别让它消耗数据库连接。

清洗是"把脏数据修到契约口径"。第 2 章模型里声明的口径(价格去币种、库存布尔化、标题去空白)在这里兑现:
from itemadapter import ItemAdapter class CleanPipeline: def process_item(self, item, spider): ad = ItemAdapter(item) ad["title"] = (ad.get("title") or "").strip() price = ad.get("price") ad["price"] = float(price.strip("£").replace(",", "")) if price else None stock_text = (ad.get("stock") or "").lower() ad["stock"] = "in stock" in stock_text return item
清洗层的纪律:只修数据格式,不猜业务含义。价格缺了就置 None 交由校验层裁决,不要自作聪明填零——填零的假数据比缺失更难排查。
校验层裁决"这件货要不要"。不合格的 Item 抛 DropItem,直接退出流水线:
from scrapy.exceptions import DropItem class ValidatePipeline: def process_item(self, item, spider): ad = ItemAdapter(item) if not ad.get("title"): raise DropItem("缺标题,作废") if ad.get("price") is None: raise DropItem("无价格,作废") return item class DedupPipeline: def open_spider(self, spider): self.seen = set() def process_item(self, item, spider): ad = ItemAdapter(item) key = ad.get("detail_url") # 主键选业务唯一键,不是自增 id if key in self.seen: raise DropItem("重复件 %s" % key) self.seen.add(key) return item
去重主键的选择是本层的灵魂:detail_url 这类业务唯一键,而不是数据库自增号。注意它与 3.1 节请求级去重的分工——那边防"同一地址下载两次",这边防"不同地址产出同一货物"(比如详情页有两个入口地址),两道闸互为补位。
落库是最贵的一层,批量是第一优化:攒一批、一次提交,吞吐提升常以十倍计:
class MySqlPipeline: def open_spider(self, spider): import pymysql self.conn = pymysql.connect(host="db.local", user="etl", password="***", database="books") self.buf = [] def process_item(self, item, spider): ad = ItemAdapter(item) self.buf.append((ad["title"], ad["price"], ad["stock"], ad["detail_url"])) if len(self.buf) >= 200: # 攒够一批再写 self.flush() return item def flush(self): with self.conn.cursor() as cur: cur.executemany( "INSERT INTO book (title, price, stock, url) VALUES (%s,%s,%s,%s)", self.buf, ) self.conn.commit() self.buf.clear() def close_spider(self, spider): self.flush() # 收工前把零头落库 self.conn.close()
运行统计里 item_scraped_count 与数据库行数对账:两者长期不一致,先查校验层的 DropItem 计数——差额通常都在那里。
💡 关键直觉:管道分层不是样式,是故障隔离——清洗崩了不影响落库,校验错了不脏库。把所有逻辑塞进一个 process_item 的管道,等于给流水线拆掉隔间。
用一件具体的"货"走一遍流水线,看清每层到底改了什么。进场时的原始 Item(Spider 侧产出):
title = ' A Light in the Attic\n' ← 带空白换行 price = '£51,77' ← 带币种,逗号是小数点 stock = 'In stock' ← 文本状态 detail_url = 'https://books.example.com/book/101'
清洗层输出:title 去空白为规范标题;price 去币种、逗号转点、转浮点 51.77;stock 布尔化为 True。校验层放行(字段齐全)。去重层查 detail_url,首件登记放行。落库层攒批写入。若同一天另一个入口地址产出同一 book/101,去重层抛 DropItem,统计里 dropped 计数加一——两道去重闸(3.1 请求级与本层 Item 级)各拦各的。对账公式由此成立:item_scraped_count = 落库行数 + dropped 计数 + 在途件数,不等式不成立时按层翻日志。
# 价格清洗的防御式写法:脏样本先过白名单 import re def clean_price(raw): if not raw: return None m = re.match(r"[^\d]*([\d.,]+)", raw.strip()) if not m: return None # 认不出就交校验层裁决,不硬猜 return float(m.group(1).replace(",", "."))
变式:多币种站点把币种信息一并提取成独立字段,而不是丢弃——清洗层的产出要保留下游做汇率换算的可能性,丢信息是最难回头的错。
货洗好了,下一节选仓储:文件、轻量库与关系库的选型与落地。