无尽冬日数据采集配置:AI开发多源管线搭建全流程指南
做 AI 开发,第一道门槛往往不在模型,而在数据。不管是做 RAG 知识库、Agent 技能训练,还是做垂直领域微调,前期的数据采集设置是否合理,直接决定了后面每一步的效率和质量。这次我们来看一套以“清源AI”为背景的 AI 开发通用数据采集配置方案,以“无尽冬日”作为大规模多源采集场景的项目代号,完整走一遍从环境准备、采集配置、任务调度到入库和接口暴露的流程。
这套教程解决的不是“装一个爬虫”这种小问题,而是三个更现实的问题:采集任务怎么管、增量怎么更新、批量任务怎么稳定跑。读完你会得到一套可以直接套用的目录结构、配置文件和 Python 实现模板,同时会重点说明采集的边界、授权要求和性能观察方法。适合正在做数据工程、Agent 知识库、模型微调数据集的同学参考。
1. 核心能力速览
先把这套采集设置定位清楚。按照清源AI 开发的通用思路,“无尽冬日采集设置”用来管理一条大规模、多源、可断点续跑的数据采集管线,而不是一个简单的下载脚本。
| 能力项 | 说明 |
|---|---|
| 项目类型 | AI 开发数据采集与预处理管线 |
| 采集场景 | 大规模多源数据,示例代号“无尽冬日” |
| 主要功能 | 多源采集、增量更新、任务调度、去重、清洗入库、批量任务、REST API |
| 运行环境 | Python 3.9+,Windows / Linux 均可 |
| 启动方式 | 命令行启动 + YAML 配置文件 |
| 数据库支持 | SQLite / PostgreSQL 均可,按实际项目替换 |
| 增量更新 | 支持基于内容哈希和更新时间的增量策略 |
| 批量任务 | 支持按批次目录批量处理 |
| 接口能力 | 预留 FastAPI 风格的 REST API |
| 合规边界 | 必须确认目标数据的授权与合法性 |
这里要说明一点:下面的命令和代码是通用模板,不是某个特定仓库的一键脚本。实际接入时,需要把项目路径、数据源地址、字段名替换成你自己的。
2. 采集设置的适用场景与使用边界
先说“无尽冬日”这种大规模采集场景适合做什么。最常见的三个方向:
第一,构建领域知识库。采集行业资讯、公开文档、论文摘要,处理后写入向量数据库,给 RAG 应用做检索源。第二,构建微调数据集。针对下游任务采集高质量文本,清洗后转换成指令格式。第三,做数据监控和定时同步。对指定来源做周期性快照,供分析和报表使用。
不适合什么场景?采集需要登录才能访问的非公开数据、购买后才能获取的付费内容、明显标记禁止转载的内容,都不建议碰。这里要特别强调三条边界:
- 采集前先确认目标网站的用户协议、robots 协议和版权声明。个人测试和商业使用是两个完全不同的合规等级,不能再拿“技术无国界”这种话当挡箭牌。
- 涉及用户隐私的数据一律不采。包括但不限于联系方式、个人身份信息、未脱敏的聊天记录。
- 如果后续要做模型训练、模型微调或者模型商用部署,数据授权必须复核。开源数据不等于可以随意商用,很多开源协议对使用场景有明确限制。
- 采集控制频率,避免对目标站点造成压力。尤其不要用高并发猛拉一个服务器,这既是不礼貌的行为,也更有可能触发封禁。
3. 环境准备与项目目录
3.1 基础环境检查
建议准备一台能长期运行的机器,CPU 4 核以上,内存 8GB 以上,磁盘根据数据量预留。训练模型时才需要 GPU;采集清洗阶段一般不用。
需要安装的基础组件:
- Python 3.9 以上版本
- Git(便于管理代码版本)
- SQLite3(Python 自带)或 PostgreSQL(可选)
- 如果需要跑向量化,安装对应的 embedding 依赖
推荐用虚拟环境隔离项目依赖,避免污染系统 Python。
# 创建项目目录 mkdir qingyuan-collect cd qingyuan-collect # 创建 Python 虚拟环境 python3 -m venv venv # 激活虚拟环境 source venv/bin/activate # Linux / macOS venv\Scripts\activate # Windows # 安装基础依赖 pip install requests beautifulsoup4 lxml pyyaml schedule psutil fastapi uvicorn如果网络环境受限,可以把 pip 源切换为国内镜像,这里不再展开。
3.2 项目目录结构
推荐按“配置、采集器、存储、任务、接口”分层组织,不要把所有逻辑堆在一个文件里。参考结构如下:
qingyuan-collect/ ├── config/ │ └── settings.yaml # 采集配置文件 ├── collectors/ │ ├── __init__.py │ ├── base.py # 采集器基类 │ └── endless_winter.py # 示例采集器 ├── processor/ │ ├── __init__.py │ ├── cleaner.py # 数据清洗 │ └── dedup.py # 去重逻辑 ├── storage/ │ ├── __init__.py │ └── database.py # 数据库写入 ├── tasks/ │ ├── __init__.py │ └── scheduler.py # 调度与批量任务 ├── api/ │ ├── __init__.py │ └── app.py # REST API 服务 ├── data/ │ ├── raw/ # 原始数据 │ ├── parsed/ # 解析后数据 │ └── logs/ # 运行日志 ├── main.py # 命令行入口 └── requirements.txt这个目录并不是强制要求,但建议保持“配置和代码分离、输入和输出分目录”的思路。批量任务跑起来之后,你会发现清晰的结构能少踩很多坑。
4. 采集设置配置文件详解
“清源AI 开发”里最值得花时间设计的就是配置文件。把采集源、采集频率、去重策略、存储路径都放进 YAML,代码只负责读配置和执行,后期维护会轻松很多。
下面的config/settings.yaml是一个通用模板:
project: name: "endless_winter" output_dir: "./data" sources: - name: "example_news" type: "html" url: "https://example.com/news" enabled: true interval_minutes: 120 # 采集间隔,单位分钟 headers: User-Agent: "Mozilla/5.0 (compatible; QingyuanBot/1.0)" - name: "example_articles" type: "api" url: "https://api.example.com/articles" enabled: false interval_minutes: 360 params: page_size: 50 sort: "latest" filter: min_text_length: 50 # 过滤掉过短的文本 max_text_length: 100000 # 过滤掉超长文本 keywords_include: [] # 必须包含的关键词,为空表示不限制 keywords_exclude: [] # 必须排除的关键词 dedup: enabled: true field: "content_hash" # 基于内容哈希去重 hash_algorithm: "md5" storage: type: "sqlite" # sqlite 或 postgresql sqlite_path: "./data/endless_winter.db" table_name: "documents" # postgresql 配置示例: # host: "127.0.0.1" # port: 5432 # user: "qingyuan" # password: "your_password" # database: "qingyuan_db" scheduler: timezone: "Asia/Shanghai" worker_threads: 2 # 并发采集线程数 retry_times: 3 # 单个任务失败重试次数 retry_delay_seconds: 30 # 重试等待时间 api: enabled: true host: "127.0.0.1" port: 8787配置字段的说明:
sources是采集源列表。每个源都有独立的开关和采集频率,不要把所有源共用同一个频率,否则高频源会被低频源拖慢,或者反过来给目标站点造成压力。filter控制采集后的初筛逻辑,先粗过滤再进清洗步骤,能省掉大量无效计算。dedup是数据质量的生命线。大型采集场景下,重复数据是最常见的问题,重复数据不处理,向量数据库和训练集都会被污染。scheduler里的worker_threads要结合机器性能设置,不是越大越快。如果你只是 4 核 CPU,设置 8 个采集线程反而会加剧资源竞争。api只在需要对外提供查询或任务触发能力时开启。
5. 核心采集逻辑与增量更新实现
配置文件写好后,接下来实现采集器。这里以基于 requests 和 BeautifulSoup 的 HTML 采集为例,演示一个可运行的采集器模板。
# collectors/base.py import hashlib import logging import time import requests from bs4 import BeautifulSoup logger = logging.getLogger(__name__) class BaseCollector: """采集器基类,统一处理请求、解析、去重前逻辑""" def __init__(self, source_config, config): self.source_config = source_config self.config = config self.session = requests.Session() self.session.headers.update(source_config.get("headers", {})) def fetch(self, url): """发送请求并返回响应文本""" resp = self.session.get(url, timeout=30) resp.raise_for_status() return resp.text def parse(self, html): """解析 HTML,返回待处理文本列表,子类需要重写""" raise NotImplementedError def content_hash(self, text): """计算文本内容哈希,用于去重""" algo = self.config.get("dedup", {}).get("hash_algorithm", "md5") h = hashlib.new(algo) h.update(text.encode("utf-8")) return h.hexdigest() def collect(self): """执行一次采集流程,返回 [(content, content_hash), ...]""" url = self.source_config["url"] html = self.fetch(url) items = self.parse(html) results = [] for item in items: content = item.get("content", "").strip() if len(content) < self.config["filter"]["min_text_length"]: continue results.append({ "source": self.source_config["name"], "content": content, "content_hash": self.content_hash(content), "collected_at": int(time.time()), }) return results具体站点解析逻辑写在子类中:
# collectors/endless_winter.py from .base import BaseCollector from bs4 import BeautifulSoup class EndlessWinterCollector(BaseCollector): """示例采集器:解析文章列表页""" def parse(self, html): soup = BeautifulSoup(html, "lxml") items = [] for article in soup.select("div.article-list > article"): title = article.select_one("h2.title") body = article.select_one("div.content") if title and body: items.append({ "content": f"{title.get_text(strip=True)}\n{body.get_text(strip=True)}" }) return items这个示例对应的是一个非常典型的列表页结构,实际站点需要根据页面结构调整 CSS 选择器。
增量更新的思路是:每次采集得到的content_hash先查数据库,如果已存在则跳过,不存在才插入。这样第二次跑的时候只处理新增内容,既不重复入库,也降低了存储压力。
# storage/database.py import sqlite3 import os class Database: """SQLite 写入与去重查询""" def __init__(self, config): self.config = config db_path = config["storage"]["sqlite_path"] os.makedirs(os.path.dirname(db_path), exist_ok=True) self.conn = sqlite3.connect(db_path) self._create_table() def _create_table(self): sql = """ CREATE TABLE IF NOT EXISTS documents ( id INTEGER PRIMARY KEY AUTOINCREMENT, source TEXT, content TEXT, content_hash TEXT UNIQUE, collected_at INTEGER ) """ self.conn.execute(sql) self.conn.commit() def exists(self, content_hash): row = self.conn.execute( "SELECT 1 FROM documents WHERE content_hash = ?", (content_hash,) ).fetchone() return row is not None def insert_many(self, items): """批量插入,遇到唯一冲突则跳过""" sql = """ INSERT OR IGNORE INTO documents (source, content, content_hash, collected_at) VALUES (?, ?, ?, ?) """ self.conn.executemany( sql, [ (item["source"], item["content"], item["content_hash"], item["collected_at"]) for item in items ], ) self.conn.commit() return self.conn.total_changes6. 批量任务与调度设计
“无尽冬日采集设置”这类场景下,批量任务不是简单地把列表循环一遍,而是要解决三个问题:定时触发、失败重试、并发控制。
下面的调度器使用schedule库并配合线程池,实现多源独立调度:
# tasks/scheduler.py import logging import threading import time from concurrent.futures import ThreadPoolExecutor import schedule from storage.database import Database logger = logging.getLogger(__name__) class TaskScheduler: def __init__(self, config, collector_map): self.config = config self.collector_map = collector_map self.storage = Database(config) self.executor = ThreadPoolExecutor( max_workers=config["scheduler"]["worker_threads"] ) def run_once(self, source_name): """执行单个采集源的一次采集""" source_config = None for s in self.config["sources"]: if s["name"] == source_name: source_config = s break if source_config is None: logger.error("source config not found: %s", source_name) return collector_cls = self.collector_map.get(source_name) if collector_cls is None: logger.error("collector not registered: %s", source_name) return collector = collector_cls(source_config, self.config) try: items = collector.collect() if not items: logger.info("no new items for source: %s", source_name) return new_items = [] for item in items: # 先做增量判断,避免写入已存在的数据 if not self.storage.exists(item["content_hash"]): new_items.append(item) if new_items: self.storage.insert_many(new_items) logger.info( "source=%s inserted=%d total_candidates=%d", source_name, len(new_items), len(items), ) except Exception: logger.exception("collect failed: %s", source_name) def run_all_once(self): """手动执行所有已启用的采集源""" for source in self.config["sources"]: if not source.get("enabled", True): continue self.executor.submit(self.run_once, source["name"]) def start_scheduler(self): """启动定时调度""" for source in self.config["sources"]: if not source.get("enabled", True): continue interval = source["interval_minutes"] schedule.every(interval).minutes.do( lambda name=source["name"]: self.executor.submit( self.run_once, name ) ) logger.info("scheduled source=%s interval=%dmin", source["name"], interval) while True: schedule.run_pending() time.sleep(1)主入口main.py做两件事:注册采集器到采集源名,然后支持“立即执行一次”和“定时循环”两种模式。
# main.py import logging import yaml from collectors.endless_winter import EndlessWinterCollector from tasks.scheduler import TaskScheduler logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s - %(message)s", ) COLLECTOR_MAP = { "example_news": EndlessWinterCollector, } def load_config(path="config/settings.yaml"): with open(path, "r", encoding="utf-8") as f: return yaml.safe_load(f) def main(): import sys config = load_config() scheduler = TaskScheduler(config, COLLECTOR_MAP) mode = sys.argv[1] if len(sys.argv) > 1 else "once" if mode == "once": scheduler.run_all_once() elif mode == "schedule": scheduler.start_scheduler() else: print("Usage: python main.py [once|schedule]") if __name__ == "__main__": main()运行方式:
# 立即跑一次全部启用的采集源 python main.py once # 按配置的间隔启动定时采集 python main.py schedule需要说明的是,schedule适合单机、数路采集源的场景。如果采集源上百个、单日数据量在百万级,建议换成 Celery 或 Arq 这样的分布式任务队列,并引入 Redis 做任务状态管理。这部分取决于你的数据规模,不要一上来就上重组件。
7. 数据清洗与入库
采集器拿到的原始文本通常不能直接用。常见问题包括:HTML 标签残留、特殊字符、重复空行、编码混乱、无意义短句。
清洗流程建议按顺序执行:
# processor/cleaner.py import re class TextCleaner: @staticmethod def remove_html_tags(text): return re.sub(r"<[^>]+>", "", text) @staticmethod def remove_urls(text): return re.sub(r"https?://\S+", "", text) @staticmethod def normalize_whitespace(text): text = re.sub(r"\s+", " ", text) return text.strip() @staticmethod def filter_by_length(text, min_len=50, max_len=100000): return min_len <= len(text) <= max_len def clean(self, text): text = self.remove_html_tags(text) text = self.remove_urls(text) text = self.normalize_whitespace(text) return text清洗时注意两个原则:
第一,清洗规则要保守。宁可少删,不要多删。比如去除 URL 时,如果文本是技术文档,URL 本身可能承载参考链接信息,直接删掉会影响语义。这种情况下更好的策略是保留一个reference_links字段,把 URL 单独提取出来。
第二,清洗后要重新进行文本长度过滤。因为清洗会缩短文本,原来刚过最短长度限制的内容,清洗后可能变得过短,需要再滤一次。
入库阶段比较简单,SQLite 的批量插入在前面已经给出了。如果数据规模增长到百万级,建议迁移到 PostgreSQL,并给content_hash字段建唯一索引,这样INSERT OR IGNORE的去重效率会好很多。
8. 接口 API 与外部接入
采集任务往往不是孤立运行的,后续要接到自己的数据处理平台或者业务系统里。预留一个轻量 API 服务很有必要。
使用 FastAPI 提供一个最小可用的接口示例,功能包括:获取当前配置、手动触发一次全量采集、查看数据库数据量:
# api/app.py from fastapi import FastAPI from pydantic import BaseModel from storage.database import Database from tasks.scheduler import TaskScheduler app = FastAPI(title="Qingyuan Collect Service") db = None scheduler = None @app.on_event("startup") def startup(): global db, scheduler from main import load_config, COLLECTOR_MAP config = load_config() db = Database(config) scheduler = TaskScheduler(config, COLLECTOR_MAP) class SourceItem(BaseModel): source: str @app.get("/health") def health(): return {"status": "ok"} @app.get("/stats") def stats(): total = db.conn.execute("SELECT COUNT(*) FROM documents").fetchone()[0] return {"total_documents": total} @app.post("/collect/run") def run_collect(): """手动触发一次全量采集""" scheduler.run_all_once() return {"status": "accepted"} @app.post("/collect/source") def run_source(item: SourceItem): """手动触发指定采集源""" scheduler.executor.submit(scheduler.run_once, item.source) return {"status": "accepted", "source": item.source}启动接口服务:
uvicorn api.app:app --host 127.0.0.1 --port 8787调用示例:
# 健康检查 curl http://127.0.0.1:8787/health # 查看数据量 curl http://127.0.0.1:8787/stats # 触发全量采集 curl -X POST http://127.0.0.1:8787/collect/run # 触发某个采集源 curl -X POST http://127.0.0.1:8787/collect/source \ -H "Content-Type: application/json" \ -d '{"source": "example_news"}'需要提醒一点:这个 API 示例没有鉴权,只适合在本地或内网使用。如果要部署到可被外部访问的环境中,必须加上 API Key 或 OAuth 认证,同时限制访问来源 IP。
9. 资源占用与性能观察
采集管线的性能观察不要只看 CPU,要重点关注四个方面:进程内存、线程池饱和度、数据库写入延迟、目标站点响应时间。
9.1 查看 CPU 和内存使用
运行中可以用系统工具快速观察:
# 查看进程占用 top -p $(pgrep -f main.py) # 更详细的进程信息 ps aux | grep main.py9.2 加入日志统计
在run_once方法中已经有基础日志。建议再增加耗时统计,判断每个采集源是否出现了慢请求:
import time start = time.time() items = collector.collect() elapsed = time.time() - start logger.info("source=%s elapsed=%.2fs items=%d", source_name, elapsed, len(items))如果某个源耗时持续增大,大概率是目标页面结构变化导致解析变慢,或者请求被限速。
9.3 显存与 GPU 观察方法
这里需要特别说明:如果采集管线中不包含模型推理,完全不需要 GPU。只有在后续做 embedding 向量化或本地模型推理时,才需要观察显存占用。
观察方法:
- 使用
nvidia-smi -l 1每 1 秒刷新一次显存占用。 - 使用
nvidia-smi --query-gpu=memory.used,utilization.gpu --format=csv输出可读的显存数据。 - 在 Python 中可以使用
pynvml获取显存信息并写入日志。
显存占用会随 batch size、文本长度、模型参数量级变化,实际占用需要以你的模型版本和推理参数为准,不能直接用别人的数字套用。
9.4 降低资源占用的技巧
- 控制并发线程数。默认 2 个采集线程对多数小规模场景足够。
- 不要每批次都打开数据库连接。持续复用连接,写入性能会稳定很多。
- 对采集到的原始 HTML 做临时落盘时,定期清理,避免磁盘空间被中间文件占满。
- 设置请求超时。所有请求都要有明确的
timeout,否则一个慢接口可能拖住整个调度线程池。
10. 常见问题与排查方法
采集管线的故障类型相对固定,下面整理一张排查表,实际运维时可以直接对照。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 启动后提示模块找不到 | 依赖未安装 | 执行pip list检查依赖 | 按 requirements.txt 重新安装依赖 |
| 配置文件读取失败 | YAML 格式错误或字段缺失 | 查看控制台报错信息,检查缩进 | 用yaml.safe_load单独加载配置文件测试 |
| 请求被拒绝或超时 | 目标站点限流、IP 被临时封禁 | 查看状态码,尝试单独请求目标 URL | 降低采集频率,加入随机延迟,使用合规代理 |
| 采集到的内容为空 | 页面结构变化、CSS 选择器失效 | 手动访问 URL 检查页面结构 | 更新解析逻辑,增加结构变化告警 |
| 数据重复入库 | 去重字段未生效或哈希算法不一致 | 检查数据库唯一索引,确认content_hash字段 | 为content_hash建唯一索引,检查哈希编码 |
| 定时任务不触发 | 时区配置错误或任务异常退出 | 查看调度日志,确认进程是否存活 | 校验时区配置,添加进程守护 |
| 数据库写入变慢 | 数据量过大、缺少索引 | 查看数据库表大小,执行查询计划 | 迁移到 PostgreSQL,按采集时间建立索引 |
| API 服务无法访问 | 端口被占用或未绑定正确地址 | 执行lsof -i:8787或netstat -ano | 更换端口或释放占用进程 |
最容易被忽略的是日志。采集任务失败后,第一件事永远是去看最近一条日志,而不是盲目重跑。建议在采集器里用logger.exception记录完整堆栈,不要只输出一行“采集失败”。
11. 最佳实践与合规建议
从“无尽冬日”这个案例延伸出去,大规模采集设置会长期运行,稳定性和合规性比单次采集量更重要。这里给几条实在的建议:
- 第一次先小规模验证。不要上来就开全量采集,先拿 100 条数据跑通整个链路,确认清洗和入库结果再放大规模。
- 保留一套最小可运行配置。所有新增采集源先在配置里
enabled: false,确保解析逻辑没问题后再开启。 - 模型文件、输入素材、输出结果分目录管理。不要让临时文件和长期数据混在一起,这样备份和清理都会更简单。
- 批量任务一定要加失败重试和日志。失败任务要能从断点继续,而不是全部从头开始。
- API 服务必须限制访问范围。默认绑定
127.0.0.1,需要对外暴露时加认证和访问控制。 - 涉及人脸、声音、版权素材数据时,必须确认授权。这一点在语音数据和视频数据采集时尤其重要,不是“只用于研究”就能免责。
- 发布或商用前做效果复核。采集来的数据入库后要抽检,不能只看数量不看质量。
另外,如果采集的数据要用于模型训练或微调,建议保留数据来源信息,包括原始 URL、采集时间、授权状态。这既是合规审计的需要,也是后续做数据溯源和错误修正的基础。
12. 总结与下一步
“清源AI 开发”中的采集设置,核心不是写一个请求库循环,而是把采集、调度、去重、清洗、入库、API 串成一条稳定可维护的管线。这篇文章以“无尽冬日”为项目代号,给出了完整的配置文件和 Python 实现模板,但实际接入时,你需要重点关注三件事:目标数据源的合法性、增量去重策略、批量任务的失败恢复流程。
建议先跑通once模式,用一个小数据集验证整个链路,确认入库数据质量后再启动定时调度。最容易踩的坑是页面结构变化导致解析器失效,所以日志和告警要比采集功能更早完善。
后续可以继续扩展的方向包括:接入向量数据库做语义检索、引入 Celery 做分布式任务队列、增加采集质量指标的可视化面板。如果采集的数据源数量很大,还可以给每种数据源单独配置解析模板,并在 Web 管理界面里动态更新,降低维护成本。
