无尽冬日数据采集配置:AI开发多源管线搭建全流程指南 做 AI 开发第一道门槛往往不在模型而在数据。不管是做 RAG 知识库、Agent 技能训练还是做垂直领域微调前期的数据采集设置是否合理直接决定了后面每一步的效率和质量。这次我们来看一套以“清源AI”为背景的 AI 开发通用数据采集配置方案以“无尽冬日”作为大规模多源采集场景的项目代号完整走一遍从环境准备、采集配置、任务调度到入库和接口暴露的流程。这套教程解决的不是“装一个爬虫”这种小问题而是三个更现实的问题采集任务怎么管、增量怎么更新、批量任务怎么稳定跑。读完你会得到一套可以直接套用的目录结构、配置文件和 Python 实现模板同时会重点说明采集的边界、授权要求和性能观察方法。适合正在做数据工程、Agent 知识库、模型微调数据集的同学参考。1. 核心能力速览先把这套采集设置定位清楚。按照清源AI 开发的通用思路“无尽冬日采集设置”用来管理一条大规模、多源、可断点续跑的数据采集管线而不是一个简单的下载脚本。能力项说明项目类型AI 开发数据采集与预处理管线采集场景大规模多源数据示例代号“无尽冬日”主要功能多源采集、增量更新、任务调度、去重、清洗入库、批量任务、REST API运行环境Python 3.9Windows / Linux 均可启动方式命令行启动 YAML 配置文件数据库支持SQLite / PostgreSQL 均可按实际项目替换增量更新支持基于内容哈希和更新时间的增量策略批量任务支持按批次目录批量处理接口能力预留 FastAPI 风格的 REST API合规边界必须确认目标数据的授权与合法性这里要说明一点下面的命令和代码是通用模板不是某个特定仓库的一键脚本。实际接入时需要把项目路径、数据源地址、字段名替换成你自己的。2. 采集设置的适用场景与使用边界先说“无尽冬日”这种大规模采集场景适合做什么。最常见的三个方向第一构建领域知识库。采集行业资讯、公开文档、论文摘要处理后写入向量数据库给 RAG 应用做检索源。第二构建微调数据集。针对下游任务采集高质量文本清洗后转换成指令格式。第三做数据监控和定时同步。对指定来源做周期性快照供分析和报表使用。不适合什么场景采集需要登录才能访问的非公开数据、购买后才能获取的付费内容、明显标记禁止转载的内容都不建议碰。这里要特别强调三条边界采集前先确认目标网站的用户协议、robots 协议和版权声明。个人测试和商业使用是两个完全不同的合规等级不能再拿“技术无国界”这种话当挡箭牌。涉及用户隐私的数据一律不采。包括但不限于联系方式、个人身份信息、未脱敏的聊天记录。如果后续要做模型训练、模型微调或者模型商用部署数据授权必须复核。开源数据不等于可以随意商用很多开源协议对使用场景有明确限制。采集控制频率避免对目标站点造成压力。尤其不要用高并发猛拉一个服务器这既是不礼貌的行为也更有可能触发封禁。3. 环境准备与项目目录3.1 基础环境检查建议准备一台能长期运行的机器CPU 4 核以上内存 8GB 以上磁盘根据数据量预留。训练模型时才需要 GPU采集清洗阶段一般不用。需要安装的基础组件Python 3.9 以上版本Git便于管理代码版本SQLite3Python 自带或 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, timeout30) 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(stripTrue)}\n{body.get_text(stripTrue)} }) 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_okTrue) 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_workersconfig[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 namesource[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( levellogging.INFO, format%(asctime)s %(levelname)s %(name)s - %(message)s, ) COLLECTOR_MAP { example_news: EndlessWinterCollector, } def load_config(pathconfig/settings.yaml): with open(path, r, encodingutf-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(rhttps?://\S, , text) staticmethod def normalize_whitespace(text): text re.sub(r\s, , text) return text.strip() staticmethod def filter_by_length(text, min_len50, max_len100000): 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(titleQingyuan 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-gpumemory.used,utilization.gpu --formatcsv输出可读的显存数据。在 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 管理界面里动态更新降低维护成本。