写入 MySQL:批量插入与去重

爬虫最脆弱的环节往往不是抓取,而是入库:单条 INSERT 一秒只能写几百行,网络抖动一次整个进程就崩,重复跑一遍又产生一堆脏数据。本章给出建表思路、参数化写入、executemany 批量提交、去重方案对照,以及断线重连与常见报错的处置。

先设计表:唯一键比代码里的去重更可靠

字段建议类型说明
idBIGINT AUTO_INCREMENT代理主键,不暴露给业务
urlVARCHAR(1024)原文链接;太长无法直接建唯一索引
url_md5CHAR(32)md5(url),用它做唯一键,索引更小更快
titleVARCHAR(512)预留冗余长度,宁可截断也不要写失败
pub_timeDATETIME统一存 UTC 或统一存本地时区,不要混存
contentMEDIUMTEXT正文;确认 max_allowed_packet 足够大
collect_timeDATETIME默认 CURRENT_TIMESTAMP,便于按批回溯补跑
CREATE TABLE news (
  id           BIGINT AUTO_INCREMENT PRIMARY KEY,
  url          VARCHAR(1024)  NOT NULL,
  url_md5      CHAR(32)       NOT NULL COMMENT 'md5(url),唯一键去重',
  title        VARCHAR(512)   NOT NULL DEFAULT '',
  content      MEDIUMTEXT,
  pub_time     DATETIME       NULL,
  collect_time DATETIME       NOT NULL DEFAULT CURRENT_TIMESTAMP,
  UNIQUE KEY uk_url_md5 (url_md5),
  KEY idx_pub_time (pub_time)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci;

唯一键选 URL 指纹而不是 title:标题会重复、会被编辑,URL 才是天然主键。若同一商品/职位要按天留多份历史,把唯一键改成 (biz_id, collect_date) 复合键。

参数化写入与 executemany

# db_writer.py —— Python 3.10+,依赖 pymysql
import hashlib, time
import pymysql

class MysqlWriter:
    """攒批 → executemany → 事务提交,失败自动重连重试"""
    INSERT_SQL = ("INSERT INTO news (url, url_md5, title, content, pub_time) "
                  "VALUES (%s, %s, %s, %s, %s) "
                  "ON DUPLICATE KEY UPDATE title=VALUES(title), content=VALUES(content)")

    def __init__(self, **kwargs):
        self.kwargs, self.conn = kwargs, None

    def connect(self):
        # autocommit=False:由我们控制事务边界,便于批量提交与回滚
        self.conn = pymysql.connect(charset="utf8mb4", autocommit=False, **self.kwargs)

    def _ensure_conn(self):
        """连接失效时重建(服务端 wait_timeout 会回收空闲连接)"""
        try:
            if self.conn is None or not self.conn.open:
                self.connect()
        except Exception:
            self.connect()

    def upsert_many(self, rows: list[dict], retry: int = 3) -> int:
        if not rows:
            return 0
        params = [(r["url"], hashlib.md5(r["url"].encode("utf-8")).hexdigest(),
                   r["title"][:512], r.get("content"), r.get("pub_time")) for r in rows]
        for attempt in range(1, retry + 1):
            try:
                self._ensure_conn()
                with self.conn.cursor() as cur:
                    cur.executemany(self.INSERT_SQL, params)   # 合成一次网络往返
                self.conn.commit()
                return len(params)
            except pymysql.err.OperationalError as exc:        # 2006/2013 连接断开
                print(f"第 {attempt} 次写入失败:{exc}")
                self.conn = None
                time.sleep(min(2 ** attempt, 10))
        raise RuntimeError(f"写入失败,共重试 {retry} 次")

if __name__ == "__main__":
    writer = MysqlWriter(host="127.0.0.1", port=3306, user="crawler",
                         password="***", database="spider")
    print(writer.upsert_many([{"url": "https://example.com/1", "title": "示例文章一"}]))

executemany 只是把多条 SQL 一次发给 MySQL,真正的吞吐来自分批大小:每批 200~1000 行最划算,再大会顶到 max_allowed_packet,再小则网络往返占比过高。

去重的几种方案

方案写法吞吐风险
先查后插SELECT 1 ... WHERE url_md5=%s 再插入低,每条多一次往返并发下仍可能重复,需唯一键兜底
唯一键 + INSERT IGNORE冲突直接忽略内容更新时旧行不更新,数据会过期
唯一键 + ON DUPLICATE KEY UPDATE冲突时更新指定列会覆盖人工修正过的字段
唯一键 + REPLACE INTO先删后插主键变化、自增 ID 膨胀
自拼多值 VALUES (...),(...)一次插上百行最高需自行转义,容易写出注入

推荐默认组合:唯一键 + ON DUPLICATE KEY UPDATE,只更新确实会变的列(标题、正文、状态),不更新人工维护的列。批量写太大则失败重试代价高、锁持有时间长,按 500 行一页提交即可:

import json

def flush_paged(writer: MysqlWriter, rows: list[dict], page: int = 500) -> int:
    """分页提交:每页一个事务,失败只影响当前页,便于定位与补偿"""
    total = 0
    for start in range(0, len(rows), page):
        chunk = rows[start:start + page]
        try:
            total += writer.upsert_many(chunk)
        except RuntimeError as exc:
            # 单页失败不中断整批,坏数据落盘待人工核对后重放
            with open("output/failed_rows.jsonl", "a", encoding="utf-8") as f:
                for row in chunk:
                    f.write(json.dumps({"url": row["url"], "error": str(exc)},
                                       ensure_ascii=False) + "\n")
    return total

「边采边入库」还是「先落文件再入库」

方案优点缺点建议
采集进程直接写库链路短、数据实时库抖动阻塞采集,坏数据直接入库数据源稳定、量不大时用
先写 JSONL 再离线入库采集与入库解耦,可重放可校验多一步,数据有延迟数据源不稳或需审计时用
写 Kafka/Redis 再消费削峰、多下游共用运维成本高多系统消费同一份数据时用

常见坑

报错/现象原因处理
Data too long for column 'title'字段长度小于实际内容写入前截断,或按 P99 长度重设字段
Packet too large单批超过 max_allowed_packet缩小批量或调大参数
Incorrect string value表/连接字符集不是 utf8mb4建表与连接都指定 utf8mb4
时间是 8 小时前的服务端与程序时区不一致统一存 UTC,展示时再转换
越写越慢没提交,形成 InnoDB 长事务每批显式 commit
重复数据仍出现唯一键建在 url(1024) 上被截断url_md5 建唯一键

合规提示:入库前确认数据来源合法——只采集公开数据,遵守 robots.txt 与站点服务条款,控制频率与并发;不采集个人隐私信息(手机号、身份证号、简历个人信息等),不绕过登录、验证码与付费墙。存储与使用需符合《个人信息保护法》《数据安全法》《著作权法》,商业再分发前取得授权,优先使用官方 API 与开放数据。

小结:把去重交给数据库唯一键,而不是应用层的「先查后插」;url_md5 做唯一键,executemany 按 200~1000 行分批,ON DUPLICATE KEY UPDATE 只更新会变的列,每批显式 commit,断线要重连重试;数据源不稳时先用 JSONL 落地再离线入库。

笔记加载中…