面向工程闭环的数据采集教学系统:requests+tenacity+Parquet实战 简介本资源是一套面向高校数据科学与大数据技术课程的教学实践源码包专为数据采集技术作业设计与实操训练打造适用于本科生课程教学、实验课指导及自学进阶。压缩包共69个文件含38.08MB涵盖9个HTML页面实现前端交互与表单展示、8个Python脚本承担数据抓取、解析、存储等核心逻辑、48张PNG/JPG图片含流程图、界面截图、可视化示例等辅助教学素材以及readme.txt说明文档和分次作业目录结构如第一次至第四次作业、2023.9.12专项任务等层次清晰、便于按教学进度分模块使用。目前已有328人学习下载。学习者可直接运行Python脚本完成真实网页爬取任务结合HTML页面理解前后端协同机制并通过图像素材直观掌握数据采集各环节的设计逻辑与输出效果配套的结构化目录与作业命名体系也极大降低了教学分发与自主学习门槛。1. 这不是“爬虫作业模板”而是面向真实工程闭环的数据采集教学设计从课堂任务到可部署脚本的完整链路你交过多少次“用 requests 抓豆瓣电影 Top250”的数据采集作业改完又改老师批注“缺乏异常处理”“没考虑反爬”“数据结构不规范”但下一次还是照着旧模板抄——这不是学生懒是课程设计本身断层了数据科学与大数据技术课教的是 Spark 和 Hadoop 架构但数据采集环节却卡在 Python 基础语法层面连 HTTP 状态码都靠 print 调试更别说对接下游 Flink 流处理或 Hive 表分区。这个标题里的“作业设计源码”本质是一套可验证、可调试、可嵌入生产 pipeline 的教学级采集系统它强制要求学生定义 schema、写 retry 逻辑、加 UA 池、存 Parquet 分区、打时间戳标签最后用 PySpark 读取并做字段统计——所有操作都在本地完成但每一步都对标企业级数据中台的采集规范。适合高校教师重构实验课、助教出题、学生做课程设计答辩也适合刚转行的数据工程师补全“从数据源头到数仓建模”的实操断点。它不教你怎么写 for 循环而教你为什么采集失败时要重试 3 次而非无限循环为什么 JSON 字段要 flatten 而不是直接塞进 DataFrame为什么 timestamp 必须用 UTC 而非本地时区。2. 用 requests fake_useragent tenacity 构建健壮采集器最小可行代码与参数精调逻辑教学场景最怕“一跑就崩”。学生写的采集脚本常因网络抖动、目标站限流、JSON 解析失败直接退出导致作业无法持续运行、数据不全、老师难评分。我们不用 Selenium太重、不用 Scrapy学习曲线陡而是用 requests 打底叠加三个轻量但关键的库fake_useragent随机 UA、tenacity可控重试、requests.adapters.HTTPAdapter连接池复用。这套组合在 90% 的静态页面采集中稳定度超 95%且代码透明、易 debug、便于学生理解每一层作用。2.1 初始化带重试策略的 Session为什么不能直接 requests.get()import requests from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type from fake_useragent import UserAgent import time # 初始化全局 UA 池避免每次请求都生成新 UA ua UserAgent(use_cache_serverFalse, fallbackMozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36) retry( stopstop_after_attempt(3), # 最多重试 3 次 waitwait_exponential(multiplier1, min1, max10), # 指数退避1s, 2s, 4s retryretry_if_exception_type((requests.exceptions.ConnectionError, requests.exceptions.Timeout, requests.exceptions.HTTPError)) ) def safe_request(url, session): headers { User-Agent: ua.random, Accept: text/html,application/xhtmlxml,application/xml;q0.9,*/*;q0.8, Accept-Language: zh-CN,zh;q0.9,en-US;q0.8,en;q0.7, Connection: keep-alive } resp session.get(url, headersheaders, timeout(3.05, 27)) # connect3.05s, read27s resp.raise_for_status() # 显式抛出 4xx/5xx 异常触发 tenacity 重试 return resp # 复用连接池避免频繁创建 socket session requests.Session() adapter requests.adapters.HTTPAdapter( pool_connections10, # 同域名最大连接数 pool_maxsize10, # 总连接池大小 max_retries0 # 关闭 requests 内置重试由 tenacity 统一控制 ) session.mount(http://, adapter) session.mount(https://, adapter)提示timeout(3.05, 27)是刻意设置的非整数值——这是 HTTP/1.1 协议推荐的“connect timeout 略大于 TCP SYN 重传间隔通常 3s”避免因底层重传导致上层误判超时。max_retries0是关键requests 自带的重试会干扰 tenacity 的指数退避逻辑必须关掉。2.2 定义采集 Schema 与结构化解析从 HTML 到 Pandas DataFrame 的强约束转换很多作业失败不是因为抓不到数据而是“抓到了但存错了”。比如商品价格字段混着货币符号和空格日期格式不统一多级嵌套 JSON 没 flatten 就直接 to_csv。我们在采集前强制定义Schema类用pydantic做字段校验和类型转换确保下游无论用 Pandas 还是 Spark 读取字段名、类型、空值策略都一致。from pydantic import BaseModel, Field, validator from datetime import datetime from typing import Optional, List class ProductItem(BaseModel): product_id: str Field(..., description商品唯一ID来自URL路径) title: str Field(..., min_length1, max_length200) price: float Field(..., ge0.01, le999999.99) # 限定价格范围过滤脏数据 rating: Optional[float] Field(None, ge0.0, le5.0) review_count: int Field(0, ge0) crawled_at: datetime Field(default_factorydatetime.utcnow) # UTC 时间戳避免时区混乱 source_url: str validator(title) def strip_title(cls, v): return v.strip().replace(\n, ).replace(\t, ) validator(price) def clean_price(cls, v): if isinstance(v, str): # 清洗 ¥199.00 或 $299 格式 import re num_str re.sub(r[^\d.], , v) return float(num_str) if num_str else 0.0 return float(v) # 解析函数示例以电商列表页为例 def parse_product_list(html_content: str) - List[ProductItem]: from bs4 import BeautifulSoup soup BeautifulSoup(html_content, lxml) items [] for li in soup.select(ul.product-list li): try: item ProductItem( product_idli.get(data-id, ), titleli.select_one(.product-title).get_text(), priceli.select_one(.price).get_text(), ratingfloat(li.select_one(.rating).get(data-score, 0)), review_countint(li.select_one(.review-count).get_text().strip(()) or 0), source_urlhttps://example.com/product/ li.get(data-id, ) ) items.append(item) except Exception as e: # 记录解析失败但不中断整体流程 print(fParse failed for item: {e}) continue return items参数说明ge/le是 Pydantic 的数值边界校验default_factorydatetime.utcnow确保每个 item 带采集时刻validator方法在字段赋值前自动清洗。这样导出的 CSV/PARQUET 文件字段类型天然兼容 Spark SQL 的DECIMAL和TIMESTAMP无需额外 cast。3. 数据落地Parquet 分区存储 元数据日志让作业可追溯、可审计、可回滚学生交作业常只给一个data.csv老师没法验证是否真跑了 1000 条、有没有跳过错误页、时间戳是否伪造。我们强制采用Parquet 分区 JSON 日志双写机制数据存raw/{source}/year2024/month06/day15/元数据起始 URL、总请求数、成功数、失败详情、耗时单独存logs/crawl_20240615.json。这不仅是工程规范更是教学评估依据——老师能一眼看出学生是否真写了分页逻辑、是否处理了 404 页面、是否对 rate limit 做了 sleep。3.1 按日期分区写入 Parquet为什么不用 CSVimport pandas as pd import os from datetime import datetime, timezone def save_to_parquet(df: pd.DataFrame, base_path: str, partition_cols: List[str] None): 将 DataFrame 写入 Parquet 分区目录 :param df: 待保存的 DataFrame :param base_path: 基础路径如 data/raw/jd :param partition_cols: 分区列名列表如 [year, month, day] # 添加分区字段 now_utc datetime.now(timezone.utc) df df.copy() df[year] now_utc.year df[month] now_utc.month df[day] now_utc.day df[crawled_at] now_utc.isoformat() # 字符串格式兼容所有引擎 # 构建分区路径 partition_path os.path.join( base_path, fyear{now_utc.year}, fmonth{now_utc.month:02d}, fday{now_utc.day:02d} ) # 确保目录存在 os.makedirs(partition_path, exist_okTrue) # 写入 Parquet使用 pyarrow 引擎支持分区 df.to_parquet( partition_path, enginepyarrow, indexFalse, compressionsnappy, use_dictionaryTrue, partition_colspartition_cols or [year, month, day] ) print(f✅ Saved {len(df)} rows to {partition_path}) # 使用示例 items parse_product_list(html_content) df pd.DataFrame([item.dict() for item in items]) save_to_parquet(df, base_pathdata/raw/jd)为什么 Parquet列式存储老师查“价格平均值”只需读 price 列比 CSV 快 5~10 倍自带 schemadf.dtypes直接映射到 Spark 的StructType避免cast(double)错误分区剪枝后续用 Spark SQL 查WHERE year2024 AND month6自动跳过其他月份目录Snappy 压缩同等数据量比 CSV 小 60%上传作业包更快。3.2 生成结构化日志记录每一次采集的“数字指纹”import json from pathlib import Path def write_crawl_log( log_dir: str, start_url: str, total_requests: int, success_count: int, failed_urls: List[str], duration_sec: float, error_summary: dict ): 写入本次采集的元数据日志 :param error_summary: {status_code: count, timeout: count, parse_error: count} log_data { crawl_id: datetime.now().strftime(%Y%m%d_%H%M%S), start_url: start_url, total_requests: total_requests, success_count: success_count, failure_count: total_requests - success_count, failed_urls: failed_urls[:10], # 只存前10个失败 URL防日志过大 duration_sec: round(duration_sec, 2), error_summary: error_summary, timestamp: datetime.now(timezone.utc).isoformat(), env: { python_version: ..join(map(str, sys.version_info[:2])), requests_version: requests.__version__, pyarrow_version: pa.__version__ if pa in globals() else N/A } } log_path Path(log_dir) / fcrawl_{log_data[crawl_id]}.json log_path.parent.mkdir(parentsTrue, exist_okTrue) with open(log_path, w, encodingutf-8) as f: json.dump(log_data, f, ensure_asciiFalse, indent2) print(f Log saved to {log_path}) # 在主采集循环后调用 start_time time.time() # ... 采集逻辑 ... end_time time.time() write_crawl_log( log_dirlogs, start_urlhttps://example.com/list?page1, total_requestslen(all_urls), success_countlen(success_items), failed_urlsfailed_list, duration_secend_time - start_time, error_summary{404: 2, timeout: 1, parse_error: 0} )教学价值老师打开logs/目录按crawl_id排序就能看到学生每次作业的执行质量曲线——第一次失败 20 次第二次降到 3 次第三次 0 失败说明他真在迭代优化。这比看代码行数或 CSV 行数更有说服力。4. 本地验证 pipeline用 PySpark 读 Parquet SQL 做数据质量检查作业不能只“跑通”还要“跑得对”。我们设计了一套本地可执行的验证 pipeline用 PySpark 读取刚生成的 Parquet 分区执行 SQL 查询检查空值率、价格异常值、时间戳合理性并生成 HTML 报告。学生交作业时除了源码和数据还必须附带report_20240615.html——这是他对自己数据质量的签字画押。4.1 用 PySpark 本地模式读取 Parquet 并执行质量规则from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, count, isnan, isnull, current_timestamp, lit from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType, TimestampType import os def validate_crawl_data(parquet_path: str, report_path: str): spark SparkSession.builder \ .appName(CrawlDataValidation) \ .master(local[*]) \ .config(spark.sql.adaptive.enabled, true) \ .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ .getOrCreate() # 定义预期 schema强制类型避免 infer 导致 int 变 string expected_schema StructType([ StructField(product_id, StringType(), False), StructField(title, StringType(), False), StructField(price, DoubleType(), False), StructField(rating, DoubleType(), True), StructField(review_count, IntegerType(), False), StructField(crawled_at, TimestampType(), False), StructField(source_url, StringType(), False), StructField(year, IntegerType(), False), StructField(month, IntegerType(), False), StructField(day, IntegerType(), False) ]) # 读取 Parquet自动识别分区 df spark.read.schema(expected_schema).parquet(parquet_path) # 数据质量规则 quality_checks [ (non_null_product_id, count(when(col(product_id) , 1)).alias(empty_product_id)), (price_range, count(when(~col(price).between(0.01, 999999.99), 1)).alias(out_of_price_range)), (crawled_at_in_past, count(when(col(crawled_at) current_timestamp(), 1)).alias(future_timestamp)), (review_count_non_negative, count(when(col(review_count) 0, 1)).alias(negative_review_count)), (title_length, count(when((col(title).isNotNull()) (col(title).length() 5), 1)).alias(short_title)) ] # 执行所有检查 results {} for check_name, expr in quality_checks: result_df df.agg(expr).collect()[0] results[check_name] result_df[0] if result_df[0] is not None else 0 # 生成报告 generate_html_report(results, parquet_path, report_path) spark.stop() def generate_html_report(check_results: dict, data_path: str, output_path: str): html_content f !DOCTYPE html htmlheadmeta charsetutf-8titleData Quality Report/title styletable{{border-collapse:collapse;width:100%}}th,td{{border:1px solid #ddd;padding:8px;text-align:left}} .fail{{background-color:#ffebee}} .pass{{background-color:#e8f5e9}}/style /headbodyh2Data Quality Report for {data_path}/h2 tabletrthCheck/ththFailed Count/ththStatus/th/tr for check, count in check_results.items(): status PASS if count 0 else FAIL bg_class pass if count 0 else fail html_content ftr class{bg_class}td{check}/tdtd{count}/tdtd{status}/td/tr html_content /table/body/html with open(output_path, w, encodingutf-8) as f: f.write(html_content) print(f Quality report saved to {output_path})关键点说明spark.sql.adaptive.*开启自适应查询优化本地小数据集也能获得合理执行计划schemaexpected_schema强制类型避免price列被 infer 成 string 导致后续计算失败col(price).between(0.01, 999999.99)是业务规则不是技术限制——它让学生意识到“数据质量 业务语义 技术约束”HTML 报告用纯 CSS 实现红绿状态无需 JS学生可直接用浏览器打开验证。4.2 一键验证脚本把验证变成python validate.py --path data/raw/jd# validate.py import argparse from datetime import datetime if __name__ __main__: parser argparse.ArgumentParser(descriptionValidate crawled Parquet data) parser.add_argument(--path, requiredTrue, helpPath to Parquet directory, e.g., data/raw/jd) parser.add_argument(--report, defaultNone, helpOutput HTML report path (default: reports/report_YYYYMMDD_HHMMSS.html)) args parser.parse_args() if not args.report: timestamp datetime.now().strftime(%Y%m%d_%H%M%S) args.report freports/report_{timestamp}.html validate_crawl_data(args.path, args.report)教学闭环学生运行python validate.py --path data/raw/jd生成reports/report_20240615_143022.html打开看到全绿 PASS才敢提交作业。老师检查报告5 秒内确认数据质量达标——这才是可量化的教学成果。5. 避坑指南学生高频翻车的 4 个黑匣子以及我的血泪经验学生写采集作业90% 的问题不是不会写代码而是踩中一些“看起来正常、运行不报错、结果却错得离谱”的黑匣子。这些坑往往藏在文档角落靠百度搜不到靠问 ChatGPT 会给出错误答案。以下是我在三年助教和企业数据中台支持中整理出的最痛的 4 个坑每一条都配真实现象、根因分析和可复制的解法。5.1 现象requests.get()返回 200但response.text是空字符串或乱码response.content却有数据原因目标网站用了gzip或brBrotli压缩但 requests 默认只解压 gzip对 br 压缩返回原始二进制而.text尝试用默认编码ISO-8859-1解码二进制导致乱码或空更隐蔽的是有些站点返回Content-Encoding: identity但实际 body 是 gziprequests 不解压。解决# 方案1强制解压所有编码推荐 import gzip import brotli # pip install brotlipy def decode_response(resp: requests.Response) - str: content resp.content encoding resp.headers.get(content-encoding, ).lower() if gzip in encoding: content gzip.decompress(content) elif br in encoding: content brotli.decompress(content) return content.decode(resp.encoding or utf-8, errorsignore) # 方案2用 httpx 替代 requests更现代默认支持 br # pip install httpx import httpx resp httpx.get(url, follow_redirectsTrue, timeout30) html resp.text # 自动解码5.2 现象采集到的价格字段全是0.0但网页明明显示¥299原因价格由 JavaScript 动态渲染requests 只拿到初始 HTML不含价格而学生没意识到需要分析 AJAX 接口或模拟 JS 执行。解决第一步用浏览器开发者工具 → Network → XHR找价格接口通常是/api/product/price?id123第二步用 requests 直接调用该 API注意 Referer、X-Requested-With 等 header第三步若 API 有加密参数如 sign则用execjs执行前端 JS 生成 sign教学场景建议换目标站避免引入复杂加密教学提示在作业要求里明确标注“本任务采集静态页面若遇动态渲染请切换至 [备用URL]”。5.3 现象Parquet 文件写入成功但用 Pandasread_parquet()读取时报ArrowInvalid: Could not convert错误原因PyArrow 版本不兼容。学生本地 PyArrow 4.x而写入时用的是 12.x高版本写的文件低版本读不了或者字段含 NaN 但 schema 定义为DoubleType(nullableFalse)PyArrow 拒绝加载。解决统一环境在requirements.txt中锁定pyarrow12.0.12023 年教育场景最稳版本写入时显式处理空值# 错误df[price].fillna(0.0) → 把 NaN 变 0但业务上 0 价不合理 # 正确保持 NaN但 schema 中 price 设为 nullableTrue schema pa.schema([ pa.field(price, pa.float64(), nullableTrue) # 注意 nullableTrue ]) df.to_parquet(path, schemaschema)5.4 现象tenacity重试了 3 次但日志里显示RetryError程序退出没捕获异常原因retry装饰器默认在重试失败后抛出tenacity.RetryError而学生没用try...except包裹调用导致整个脚本崩溃。解决# ✅ 正确用法必须捕获 RetryError try: resp safe_request(url, session) except tenacity.RetryError as e: print(f❌ Failed after 3 retries for {url}: {e.last_attempt.exception()}) failed_urls.append(url) continue # 跳过当前 URL继续下一个 # ❌ 错误用法常见 # resp safe_request(url, session) # 一旦重试失败直接 exit血泪经验这些坑我当年也踩过。第一次遇到br压缩debug 了 6 小时最后发现 requests 文档里有一行小字“Brotli support requires installingbrotlipy”。所以现在我带学生第一节课就发一份《采集黑匣子清单》让他们打印贴在显示器边框上——不是为了背而是建立“遇到诡异现象先查这四条”的条件反射。6. 进阶技巧用 Docker Compose 封装采集环境实现“一键复现、跨平台交付”学生交作业最大的痛点是“我在 Windows 跑通了老师用 Mac 打不开 Parquet”“我装了 PyArrow 12同学装了 14互相读不了”。解决方案不是让大家统一环境而是用 Docker 把整个采集 runtime 打包成镜像Python 3.9、PyArrow 12.0.1、requests、tenacity、fake-useragent 全部预装连spark-submit本地模式都配置好。学生只需写docker-compose up就能在任何系统上启动完全一致的环境。6.1 Dockerfile精简镜像只装必要依赖# Dockerfile FROM python:3.9-slim # 设置时区避免日志时间错乱 ENV TZUTC RUN ln -snf /usr/share/zoneinfo/$TZ /etc/localtime echo $TZ /etc/timezone # 安装系统依赖PyArrow 编译需要 RUN apt-get update apt-get install -y \ build-essential \ libglib2.0-0 \ rm -rf /var/lib/apt/lists/* # 创建工作目录 WORKDIR /app # 复制依赖文件分离依赖安装利用 Docker cache COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt # 复制源码 COPY . . # 暴露日志和数据目录方便挂载宿主机 VOLUME [/app/data, /app/logs, /app/reports] # 默认命令 CMD [python, crawl_main.py]requirements.txt内容requests2.31.0 tenacity8.2.3 fake-useragent1.4.0 beautifulsoup44.12.2 pandas1.5.3 pyarrow12.0.1 pyspark3.4.1为什么选 slim 镜像python:3.9-slim镜像仅 120MB比python:3.9900MB小 7 倍学生下载快build-essential是 PyArrow 编译必需但slim镜像默认不带必须显式安装pyspark3.4.1是目前与 PyArrow 12 兼容最稳的版本Spark 3.5 对 Arrow 14 有强依赖。6.2 docker-compose.yml一键启动采集 验证 pipeline# docker-compose.yml version: 3.8 services: crawler: build: . volumes: - ./data:/app/data - ./logs:/app/logs - ./reports:/app/reports - ./config:/app/config:ro # 配置文件只读 environment: - PYTHONUNBUFFERED1 - SPARK_LOCAL_IP127.0.0.1 command: python crawl_main.py --source jd --pages 5 validator: build: . volumes: - ./data:/app/data - ./reports:/app/reports environment: - PYTHONUNBUFFERED1 command: python validate.py --path /app/data/raw/jd --report /app/reports/report_latest.html depends_on: - crawler # 验证服务在 crawler 完成后启动 restart: no # 可选暴露 Jupyter Lab让学生在线调试 jupyter: image: jupyter/scipy-notebook:lab-4.0.2 volumes: - ./notebooks:/home/jovyan/work - ./data:/home/jovyan/data:ro - ./logs:/home/jovyan/logs:ro ports: - 8888:8888 environment: - JUPYTER_TOKENteach2024交付效果学生交作业时不再只交.py和.csv而是交一个project.zip里面包含Dockerfiledocker-compose.ymlcrawl_main.py含采集逻辑validate.py含验证逻辑requirements.txtREADME.md含docker-compose up --build命令老师解压后终端输入docker-compose up --build等待 2 分钟自动完成采集 → 存 Parquet → 生成日志 → 运行验证 → 输出 HTML 报告。全程无需装 Python、无需配环境、无需担心版本冲突——这就是教学级工程化的底线。6.3 教学延伸如何把 Docker 环境迁移到学校机房服务器很多高校机房禁用 Docker Desktop但允许 Linux 服务器装 Docker Engine。这时只需两步在服务器上安装 Docker CEUbuntu 示例curl -fsSL https://get.docker.com | sudo bash sudo usermod -aG docker $USER # 加入 docker 组修改docker-compose.yml去掉jupyter服务机房通常不开放 8888 端口并用--network host提升网络性能services: crawler: # ... 其他配置不变 network_mode: host # 直接用宿主机网络避免 NAT 延迟我带过的 3 届学生从没一个人因为“环境配不起来”而放弃课程设计。他们可能 still 不懂 Spark Shuffle 原理但已经能自信地说“老师我的采集 pipeline 在 Docker 里跑通了您 pull 下来就能验证。” 这种确定性比教会一百个算法更重要。希望帮到你。本文还有配套的精品资源点击获取