Skills MCP Model 博客 提交 Skills
登录 注册

DeepSeek 数据工程教程

从数据采集到质量评估,完整掌握大模型数据工程流水线。涵盖 SFT 数据集构建、RLHF 偏好数据准备、数据清洗去重、数据增强、大规模数据处理等核心环节,附 Python 实战代码。

开始学习

数据是 LLM 的灵魂

在大模型训练中,数据质量直接决定了模型能力的上限。高质量的数据集能让小模型超越大模型,而低质量的数据会让万亿参数模型表现平庸。数据工程覆盖了从原始数据采集到最终训练数据交付的全链路。

数据工程概述

理解数据在 LLM 训练中的核心地位,掌握数据质量与数据数量的权衡关系,建立完整的数据工程流水线全景认知。

数据在 LLM 训练中的核心地位

在 LLM 训练的三要素(算法、算力、数据)中,数据往往被低估,但它的重要性远超其他两者。一个广为流传的观点是:模型架构的边际收益递减,而数据质量的边际收益递增。以下是三个相互关联的核心事实:

  • 数据质量决定模型能力上限:即使使用最先进的模型架构(如 MoE),如果训练数据质量差,模型也无法产出高质量回答。DeepSeek-V3 和 R1 的强大能力,很大程度上归功于精心构建的训练数据。
  • 数据多样性决定泛化能力:单一领域的数据会让模型过拟合,而多样化的数据能让模型在未见过的任务上表现良好。DeepSeek 的数据混合策略覆盖了数学、代码、推理、对话、创意写作等多个领域。
  • 数据规模决定知识边界:模型的知识范围不会超出训练数据覆盖的范围。想让模型懂医学,就必须有高质量的医学数据;想让模型会编程,就必须有足够的代码数据。

数据质量 vs 数据数量

在有限的计算资源下,数据质量比数据数量更重要。以下是两者的权衡分析:

维度 追求数量 追求质量
训练效率 训练时间长,收敛慢 训练更快,收敛更稳定
模型表现 噪声多,输出不稳定 输出准确,幻觉少
成本 GPU 成本高,周期长 数据清洗成本高,但总成本更低
典型策略 爬取全网数据,粗过滤 精选数据源,多重过滤,人工标注

DeepSeek 团队在实践中发现:用 1/10 的高质量数据训练的模型,在指令遵循和推理能力上往往优于用全量粗数据训练的模型。这也解释了为什么数据工程是 LLM 开发中最关键的环节。

数据工程流水线全景

一个完整的数据工程流水线包含以下阶段:

  1. 数据采集:从公开数据集、网络爬取、API 调用、合成生成等渠道获取原始数据
  2. 数据清洗:去除重复、过滤低质量、处理缺失值、过滤敏感信息
  3. 数据标注:为 SFT、RLHF、DPO 等不同训练阶段构建不同格式的数据
  4. 数据增强:通过 Self-Instruct、Evol-Instruct、反向翻译等技术扩充数据
  5. 质量评估:从多样性、困难度、指令复杂度等维度评估数据质量
  6. 数据混合:按比例混合不同领域的数据,形成最终的训练数据集
  7. 版本管理:使用 DVC 等工具管理数据版本,确保可复现

数据采集与来源

了解 LLM 训练数据的主要来源,包括公开数据集、网络爬取、合成数据生成,以及数据合规注意事项。

常用公开数据集

数据集名称 规模 类型 适用场景
Alpaca 52K SFT 指令数据 指令微调入门
ShareGPT 90K 多轮对话 对话能力训练
UltraChat 1.5M 多轮对话 大规模对话训练
OpenOrca 4M SFT 指令数据 大规模指令微调
CodeAlpaca 20K 代码生成 代码能力训练
MathInstruct 260K 数学推理 数学能力训练

从 Hugging Face 加载公开数据集

from datasets import load_dataset # 加载 Alpaca 数据集 alpaca = load_dataset("tatsu-lab/alpaca") print(f"Alpaca 训练集: {len(alpaca['train'])} 条") print(f"示例: {alpaca['train'][0]}") # 加载 ShareGPT 对话数据 sharegpt = load_dataset("anon8231489123/ShareGPT_Vicuna_unfiltered", data_files="ShareGPT_V3_unfiltered_cleaned_split.json") print(f"ShareGPT: {len(sharegpt['train'])} 条对话") # 加载 OpenOrca 数据(按需加载子集) orca = load_dataset("Open-Orca/OpenOrca", split="train[:50000]") print(f"OpenOrca 子集: {len(orca)} 条")

网络爬取与数据采集

对于特定领域的数据,网络爬取是重要的补充手段。以下是一个面向文档和教程的爬虫示例:

import requests from bs4 import BeautifulSoup import time from urllib.parse import urljoin, urlparse def crawl_documentation(base_url, max_pages=100, delay=1.0): """爬取文档类网站,提取正文内容""" visited = set() to_visit = [base_url] results = [] while to_visit and len(visited) < max_pages: url = to_visit.pop(0) if url in visited: continue visited.add(url) try: resp = requests.get(url, timeout=10, headers={ "User-Agent": "Mozilla/5.0 (compatible; DataBot/1.0)" }) resp.raise_for_status() soup = BeautifulSoup(resp.text, "html.parser") # 移除脚本和样式标签 for tag in soup(["script", "style", "nav", "footer"]): tag.decompose() text = soup.get_text(separator=" ", strip=True) if len(text) > 200: results.append({ "url": url, "title": soup.title.string if soup.title else "", "content": text, }) # 发现新链接 for link in soup.find_all("a", href=True): href = urljoin(url, link["href"]) if urlparse(href).netloc == urlparse(base_url).netloc: if href not in visited: to_visit.append(href) time.sleep(delay) except Exception as e: print(f"爬取失败 {url}: {e}") return results # 使用示例 docs = crawl_documentation("https://docs.python.org/3/", max_pages=50) print(f"爬取了 {len(docs)} 个页面")

合成数据生成

当公开数据不足以覆盖特定领域时,可以使用合成数据生成(Synthetic Data Generation)技术。通过 DeepSeek 等强模型生成高质量的训练数据:

from openai import OpenAI client = OpenAI( api_key="sk-your-api-key", base_url="https://api.deepseek.com/v1", ) def generate_synthetic_data(topic, num_samples=10): """使用 DeepSeek 生成特定领域的合成数据""" prompt = f"""请生成 {num_samples} 条关于「{topic}」的高质量指令数据。 每条数据包含: - instruction: 清晰、具体的指令或问题 - input: 补充上下文(可为空字符串) - output: 专业、准确、详尽的回答 输出格式为 JSON 数组。""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0.8, max_tokens=4096, ) return response.choices[0].message.content # 生成 Python 机器学习相关数据 synthetic = generate_synthetic_data("Python 机器学习", num_samples=10) print(synthetic)

数据合规注意事项

在使用公开数据集和网络爬取数据时,务必检查数据许可协议。部分数据集(如 ShareGPT)有特定的使用限制。网络爬取应遵守 robots.txt 协议,控制爬取频率,避免对目标网站造成负担。涉及个人隐私信息的数据需要脱敏处理。

数据清洗与去重

数据清洗是数据工程中最耗时但最重要的环节。高质量的数据清洗能显著提升模型训练效果,包括质量过滤、去重和敏感信息过滤。

质量过滤管道

一条完整的质量过滤管道包含多个维度的过滤:

import re from typing import List, Dict class DataQualityFilter: """多维度数据质量过滤器""" def __init__(self, min_length=20, max_length=8000, min_words=5): self.min_length = min_length self.max_length = max_length self.min_words = min_words def filter_by_length(self, text: str) -> bool: """长度过滤:太短或太长的文本可能是噪声""" return self.min_length <= len(text) <= self.max_length def filter_by_word_count(self, text: str) -> bool: """词数过滤:过滤过于碎片化的文本""" words = text.split() return len(words) >= self.min_words def filter_repetitive(self, text: str, threshold=0.3) -> bool: """重复内容过滤:检测重复 n-gram 比例""" if len(text) < 100: return True words = text.split() unigrams = len(set(words)) if unigrams / len(words) < threshold: return False # 重复率过高 return True def filter_special_chars(self, text: str, max_ratio=0.3) -> bool: """特殊字符过滤:过滤乱码或非正常文本""" special = len(re.findall(r'[^\w\s\u4e00-\u9fff.,;:!?()\-\+\=]', text)) return special / max(len(text), 1) < max_ratio def filter_empty_output(self, sample: Dict) -> bool: """空输出过滤:output 过短或仅包含占位符""" output = sample.get("output", "") placeholder_patterns = [ r'^(sorry|unfortunately|i cannot|as an ai)', r'^(抱歉|对不起|作为.*AI|我无法)', ] for pattern in placeholder_patterns: if re.match(pattern, output.strip().lower()): return False return len(output.strip()) > 20 def apply_all(self, samples: List[Dict]) -> List[Dict]: """应用所有过滤器""" filtered = [] for s in samples: text = s.get("instruction", "") + " " + s.get("output", "") if all([ self.filter_by_length(text), self.filter_by_word_count(text), self.filter_repetitive(text), self.filter_special_chars(text), self.filter_empty_output(s), ]): filtered.append(s) print(f"过滤前: {len(samples)} 条, 过滤后: {len(filtered)} 条") return filtered

MinHash LSH 去重

MinHash + LSH(Locality-Sensitive Hashing)是业界标准的大规模去重方案,能高效检测近似重复的文档:

from datasketch import MinHash, MinHashLSH import re def tokenize_chinese(text): """中文分词:按字符 2-gram 切分""" # 简化版分词,生产环境建议使用 jieba 或 pkuseg text = re.sub(r'[^\u4e00-\u9fff\w]', ' ', text.lower()) # 2-gram 分词 return [text[i:i+2] for i in range(len(text)-1)] def deduplicate_with_minhash(samples, num_perm=128, threshold=0.8): """使用 MinHash LSH 进行近似去重""" lsh = MinHashLSH(threshold=threshold, num_perm=num_perm) unique_samples = [] for idx, sample in enumerate(samples): text = sample.get("instruction", "") + " " + sample.get("output", "") tokens = tokenize_chinese(text) if len(tokens) < 10: continue m = MinHash(num_perm=num_perm) for token in tokens: m.update(token.encode("utf-8")) # 查询是否已存在近似重复 if len(lsh.query(m)) == 0: lsh.insert(str(idx), m) unique_samples.append(sample) print(f"去重前: {len(samples)} 条, 去重后: {len(unique_samples)} 条") return unique_samples

语义去重

MinHash 适合字面级去重,但对于语义相似但表述不同的数据,需要使用嵌入向量进行语义去重:

from sentence_transformers import SentenceTransformer import numpy as np from sklearn.metrics.pairwise import cosine_similarity def semantic_deduplication(samples, model_name="BAAI/bge-small-zh-v1.5", threshold=0.95, batch_size=256): """基于语义嵌入的去重""" model = SentenceTransformer(model_name) # 提取所有文本 texts = [s.get("instruction", "") for s in samples] # 批量编码 embeddings = model.encode(texts, batch_size=batch_size, show_progress_bar=True) # 计算相似度矩阵,标记重复 unique_indices = [] seen = set() for i in range(len(embeddings)): if i in seen: continue unique_indices.append(i) # 找到与当前样本高度相似的样本 sims = cosine_similarity([embeddings[i]], embeddings[i+1:])[0] duplicates = np.where(sims > threshold)[0] + i + 1 for d in duplicates: seen.add(int(d)) print(f"语义去重: {len(samples)} -> {len(unique_indices)} 条") return [samples[i] for i in unique_indices]

敏感信息过滤

在训练数据中,必须过滤掉个人隐私信息和敏感内容:

import re class SensitiveContentFilter: """敏感信息过滤器""" # 常见敏感信息正则模式 PATTERNS = { "email": r'[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}', "phone_cn": r'1[3-9]\d{9}', "id_card": r'\d{17}[\dXx]', "ip_address": r'\b(?:\d{1,3}\.){3}\d{1,3}\b', "url": r'https?://[^\s<>"{}|\\^`\[\]]+', "api_key": r'(?:sk|api[_-]?key|token)[=:]\s*[\w-]+', } def has_sensitive(self, text: str) -> bool: """检测是否包含敏感信息""" for name, pattern in self.PATTERNS.items(): if re.search(pattern, text): return True return False def mask_sensitive(self, text: str) -> str: """脱敏处理:替换敏感信息为占位符""" masked = text replacements = { "email": "[EMAIL]", "phone_cn": "[PHONE]", "id_card": "[ID_CARD]", "ip_address": "[IP]", "api_key": "[API_KEY]", } for name, replacement in replacements.items(): if name in self.PATTERNS: masked = re.sub(self.PATTERNS[name], replacement, masked) return masked # 使用示例 filter = SensitiveContentFilter() text = "请联系 admin@example.com 或拨打 13800138000" print(f"包含敏感信息: {filter.has_sensitive(text)}") # True print(f"脱敏后: {filter.mask_sensitive(text)}")

数据格式与标注

不同训练阶段需要不同的数据格式。SFT 使用 Instruction-Input-Output 格式,ChatML 用于对话场景,RLHF/DPO 需要偏好比较数据。理解这些格式是构建高质量数据集的基础。

SFT 数据格式(Instruction-Input-Output)

SFT(Supervised Fine-Tuning)数据是最基础的训练数据格式。每条数据包含指令、可选输入和期望输出:

{ "instruction": "用 Python 实现快速排序算法", "input": "", "output": "def quicksort(arr):\n if len(arr) <= 1:\n return arr\n pivot = arr[len(arr) // 2]\n left = [x for x in arr if x < pivot]\n middle = [x for x in arr if x == pivot]\n right = [x for x in arr if x > pivot]\n return quicksort(left) + middle + quicksort(right)" }

ChatML 对话格式

ChatML(Chat Markup Language)是 OpenAI 定义的对话格式,也广泛用于 SFT 训练。它将多轮对话结构化为标准格式:

def format_chatml(conversations): """将对话列表转换为 ChatML 格式字符串""" formatted = "" for turn in conversations: role = turn["role"] content = turn["content"] if role == "system": formatted += f"<|im_start|>system\n{content}<|im_end|>\n" elif role == "user": formatted += f"<|im_start|>user\n{content}<|im_end|>\n" elif role == "assistant": formatted += f"<|im_start|>assistant\n{content}<|im_end|>\n" return formatted.strip() # 示例对话 conversation = [ {"role": "system", "content": "你是一个专业的 Python 编程助手。"}, {"role": "user", "content": "如何读取 CSV 文件?"}, {"role": "assistant", "content": "使用 pandas 可以轻松读取 CSV 文件:\nimport pandas as pd\ndf = pd.read_csv('data.csv')"}, ] chatml = format_chatml(conversation) print(chatml)

RLHF 偏好数据格式

RLHF(Reinforcement Learning from Human Feedback)和 DPO(Direct Preference Optimization)需要偏好比较数据,即同一个 prompt 下,标注哪个回答更好:

# DPO 数据格式 dpo_sample = { "prompt": "解释什么是机器学习", "chosen": "机器学习是人工智能的一个分支,它使计算机能够从数据中学习模式,而无需显式编程。常见方法包括监督学习、无监督学习和强化学习。", "rejected": "机器学习就是让机器学会东西。", } # RLHF 比较数据格式 rlhf_comparison = { "prompt": "写一首关于春天的诗", "responses": [ {"text": "春风拂面柳如烟...", "score": 4.5}, {"text": "春天来了,花开了。", "score": 1.0}, ] }

数据格式选择建议

基础 SFT 训练使用 Instruction-Input-Output 格式即可;多轮对话场景推荐 ChatML 格式;如果需要进行 RLHF 或 DPO 训练,必须准备偏好对比数据。DeepSeek 系列模型同时支持 Alpaca 格式和 ChatML 格式。

数据增强技术

当已有数据不足时,数据增强技术可以帮你从少量种子数据中生成大量高质量训练数据。Self-Instruct 和 Evol-Instruct 是两种最主流的方法。

Self-Instruct 自生成方法

Self-Instruct 的核心思想是用强模型从种子任务中自动生成新的指令-输出对:

import json import random from openai import OpenAI client = OpenAI( api_key="sk-your-api-key", base_url="https://api.deepseek.com/v1", ) def generate_instructions(seed_tasks, num_to_generate=50): """基于种子任务生成新指令""" # 随机采样种子任务作为上下文 seed_context = random.sample(seed_tasks, min(8, len(seed_tasks))) seed_text = "\n".join([f"- {t['instruction']}" for t in seed_context]) prompt = f"""你是一个数据标注专家。请基于以下种子任务,生成 {num_to_generate} 个新的、多样化的指令。 种子任务示例: {seed_text} 要求: 1. 新指令应覆盖不同领域和难度级别 2. 指令应清晰、具体、可执行 3. 避免与种子任务重复 4. 输出格式为 JSON 数组,每项包含 instruction 字段 请只输出 JSON 数组,不要包含其他文字。""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0.9, max_tokens=4096, ) try: new_instructions = json.loads(response.choices[0].message.content) return new_instructions except json.JSONDecodeError: return [] def generate_output(instruction): """为给定指令生成高质量输出""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": instruction}], temperature=0.3, max_tokens=2048, ) return response.choices[0].message.content # 使用示例 seed_tasks = [ {"instruction": "用 Python 实现冒泡排序"}, {"instruction": "解释什么是 RESTful API"}, {"instruction": "写一个计算斐波那契数列的函数"}, ] new_instructions = generate_instructions(seed_tasks, num_to_generate=10) print(f"生成了 {len(new_instructions)} 条新指令")

Evol-Instruct 进化生成

Evol-Instruct 通过逐步增加指令的复杂度来生成更具挑战性的数据。DeepSeek 团队在训练中大量使用了类似的进化策略:

def evolve_instruction(instruction, evolution_type="deepen"): """进化指令:增加深度、广度或复杂度""" evolution_prompts = { "deepen": """将以下指令重写得更深入、更专业。增加对底层原理、实现细节或边缘情况的讨论。 原指令:{instruction} 深化后的指令:""", "broaden": """将以下指令重写得更广泛,要求覆盖多个相关子话题或不同场景。 原指令:{instruction} 扩展后的指令:""", "increase_reasoning": """将以下指令改写为需要多步推理才能完成的复杂任务。 原指令:{instruction} 增加推理要求后的指令:""", } prompt = evolution_prompts.get(evolution_type, evolution_prompts["deepen"]) prompt = prompt.format(instruction=instruction) response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0.7, max_tokens=1024, ) return response.choices[0].message.content # 进化示例 original = "解释什么是数据库索引" deepened = evolve_instruction(original, "deepen") print(f"原指令: {original}") print(f"深化后: {deepened}")

反向翻译数据增强

反向翻译(Back Translation)是多语言数据增强的经典方法,通过翻译-回译生成语义等价但表达不同的数据:

def back_translate(text, source_lang="中文", pivot_lang="英文"): """反向翻译:中文 -> 英文 -> 中文,生成语义等价的变体""" # 第一步:翻译为目标语言 translate_prompt = f"请将以下{source_lang}文本翻译为{pivot_lang},只输出翻译结果:\n{text}" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": translate_prompt}], temperature=0.3, ) translated = response.choices[0].message.content # 第二步:翻译回源语言 back_prompt = f"请将以下{pivot_lang}文本翻译为{source_lang},只输出翻译结果:\n{translated}" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": back_prompt}], temperature=0.3, ) return response.choices[0].message.content # 使用示例 original = "深度学习是机器学习的一个子领域,它使用多层神经网络来学习数据的表示。" augmented = back_translate(original) print(f"原文: {original}") print(f"增强: {augmented}")

DeepSeek 特定数据准备

DeepSeek-R1 的推理数据有特殊的格式要求,包含 thinking 标签和 CoT(Chain of Thought)推理链。本章详细介绍如何为 DeepSeek 模型准备数据。

DeepSeek-R1 推理数据格式

DeepSeek-R1 的关键创新在于训练数据中包含了显式的思考过程(thinking/reasoning)。以下是标准的 R1 推理数据格式:

# DeepSeek-R1 推理数据格式 r1_format_sample = { "messages": [ { "role": "user", "content": "一个长方体的长宽高分别为 3cm、4cm、5cm,求它的表面积和体积。" }, { "role": "assistant", "content": "<think>这是一个长方体表面积和体积的计算问题。\n表面积 = 2*(长*宽 + 长*高 + 宽*高)\n体积 = 长*宽*高\n\n代入数据:\n长=3, 宽=4, 高=5\n表面积 = 2*(3*4 + 3*5 + 4*5) = 2*(12+15+20) = 2*47 = 94\n体积 = 3*4*5 = 60\n\n验证:表面积和体积单位不同,但数值计算正确。</think>\n\n长方体的表面积是 94 平方厘米,体积是 60 立方厘米。" } ] }

CoT 数据构建

构建 Chain of Thought 数据的关键是让模型学会"先思考,再回答"。以下代码自动为数学问题生成 CoT 推理链:

def generate_cot_data(question, answer): """为问题和答案生成 CoT 推理过程""" prompt = f"""请为以下数学问题生成详细的逐步推理过程。 问题:{question} 正确答案:{answer} 请按以下格式输出: <think> (详细的逐步推理,包括使用的公式、中间步骤、验证过程) </think> (最终答案,简洁明了) 请确保推理过程清晰、完整、可验证。""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0.3, max_tokens=4096, ) return response.choices[0].message.content # 为数学题生成 CoT cot_result = generate_cot_data( question="若 x^2 - 5x + 6 = 0,求 x 的值。", answer="x = 2 或 x = 3" ) print(cot_result)

代码数据准备

代码数据是 DeepSeek 训练数据的重要组成部分。优质的代码数据应包含问题描述、代码实现和注释:

def format_code_data(problem, code, language="python", difficulty="medium"): """格式化代码训练数据""" return { "messages": [ { "role": "system", "content": f"你是一个专业的{language.upper()}编程助手。请提供清晰、高效的代码实现。" }, { "role": "user", "content": problem }, { "role": "assistant", "content": code } ], "metadata": { "language": language, "difficulty": difficulty, "type": "code_generation", } } # 示例 code_sample = format_code_data( problem="实现一个 LRU 缓存(最近最少使用),支持 get 和 put 操作,时间复杂度 O(1)。", code="""from collections import OrderedDict class LRUCache: def __init__(self, capacity: int): self.capacity = capacity self.cache = OrderedDict() def get(self, key: int) -> int: if key not in self.cache: return -1 self.cache.move_to_end(key) return self.cache[key] def put(self, key: int, value: int) -> None: if key in self.cache: self.cache.move_to_end(key) self.cache[key] = value if len(self.cache) > self.capacity: self.cache.popitem(last=False)""", difficulty="medium" )

数学数据准备

数学推理数据需要包含公式、推导步骤和最终答案。LaTeX 格式是数学表达的标准:

math_sample = { "messages": [ { "role": "user", "content": "求极限 \\lim_{x \\to 0} \\frac{\\sin x}{x}" }, { "role": "assistant", "content": "<think>这是一个经典极限问题。\n\n方法一:使用洛必达法则\n当 x→0 时,分子 sin x → 0,分母 x → 0,满足 0/0 型不定式。\n对分子分母分别求导:\nlim(sin x / x) = lim(cos x / 1) = cos(0) = 1\n\n方法二:使用夹逼定理\n对于 0 < |x| < pi/2,有 cos x < sin x / x < 1\n当 x→0 时,cos x → 1,因此极限为 1。\n\n验证:这个结果是正确的。</think>\n\n\\lim_{x \\to 0} \\frac{\\sin x}{x} = 1" } ] }

数据质量评估

数据质量评估是数据工程的最后一道防线。从多样性、困难度、指令复杂度、自动评分等多个维度评估数据集质量,确保训练数据符合预期标准。

多样性评估

多样性评估确保数据集覆盖了足够广泛的领域和任务类型:

from sklearn.feature_extraction.text import TfidfVectorizer from sklearn.cluster import KMeans import numpy as np def evaluate_diversity(samples, n_clusters=10): """评估数据集的多样性""" # 提取指令文本 instructions = [s.get("instruction", "") for s in samples] # TF-IDF 向量化 vectorizer = TfidfVectorizer(max_features=5000, ngram_range=(1, 2)) vectors = vectorizer.fit_transform(instructions) # 聚类分析 kmeans = KMeans(n_clusters=n_clusters, random_state=42) labels = kmeans.fit_predict(vectors.toarray()) # 统计每个聚类的样本数 cluster_counts = np.bincount(labels) cluster_ratio = cluster_counts / len(samples) # 计算多样性分数(基于聚类分布的熵) entropy = -np.sum(cluster_ratio * np.log(cluster_ratio + 1e-10)) max_entropy = np.log(n_clusters) diversity_score = entropy / max_entropy print(f"聚类数量: {n_clusters}") print(f"各聚类样本数: {cluster_counts}") print(f"多样性分数: {diversity_score:.4f} (1.0 = 完美均匀分布)") return { "diversity_score": diversity_score, "cluster_counts": cluster_counts.tolist(), "cluster_ratio": cluster_ratio.tolist(), }

困难度评估

评估数据集中每条指令的难度,确保数据难度分布合理:

def evaluate_difficulty(instruction): """使用 LLM 评估指令的困难度""" prompt = f"""请评估以下指令的困难度,返回 1-5 的评分(1=非常简单,5=非常困难)。 评估标准: - 1分:简单的事实性问题,无需推理 - 2分:需要基本理解或简单操作 - 3分:需要多步思考或综合知识 - 4分:需要高级推理或专业知识 - 5分:需要深度推理、创造性或专家级知识 指令:{instruction} 请只输出分数(1-5 的整数)。""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0, max_tokens=10, ) try: return int(response.choices[0].message.content.strip()) except ValueError: return 3 def analyze_difficulty_distribution(samples): """分析数据集的难度分布""" from collections import Counter difficulties = [] for s in samples[:100]: # 采样评估 score = evaluate_difficulty(s.get("instruction", "")) difficulties.append(score) dist = Counter(difficulties) for level in range(1, 6): print(f"难度 {level}: {dist.get(level, 0)} 条 ({dist.get(level,0)/len(difficulties)*100:.1f}%)") print(f"\n平均难度: {np.mean(difficulties):.2f}") return difficulties

指令复杂度评分

指令复杂度评分从多个维度评估指令的质量:

评分维度 说明 评分方法
清晰度 指令是否明确、无歧义 LLM 评分 1-5
完整性 指令是否包含足够的上下文 信息密度计算
可执行性 指令是否可以被完成 LLM 可执行性判断
创造性 指令是否鼓励创造性思维 LLM 评分 1-5

自动质量打分

使用 LLM 作为评判器,对数据质量进行自动评分:

def auto_quality_score(sample): """使用 LLM 自动评估单条数据质量""" instruction = sample.get("instruction", "") output = sample.get("output", "") prompt = f"""请评估以下问答对的质量,从 1-10 分打分。 评估维度: 1. 指令是否清晰明确(1-3分) 2. 答案是否准确完整(1-3分) 3. 答案是否专业、有帮助(1-2分) 4. 答案长度是否适中(1-2分) 指令:{instruction} 答案:{output} 请输出 JSON 格式: {{"total_score": X, "reasons": ["原因1", "原因2"]}}""" response = client.chat.completions.create( model="deepseek-chat", messages=[{"role": "user", "content": prompt}], temperature=0, max_tokens=512, ) try: return json.loads(response.choices[0].message.content) except json.JSONDecodeError: return {"total_score": 5, "reasons": ["无法解析评分"]} # 使用示例 sample = { "instruction": "用 Python 实现二分查找算法,要求处理边界条件", "output": "def binary_search(arr, target):\n left, right = 0, len(arr)-1\n while left <= right:\n mid = (left+right)//2\n if arr[mid] == target:\n return mid\n elif arr[mid] < target:\n left = mid+1\n else:\n right = mid-1\n return -1" } score = auto_quality_score(sample) print(f"质量评分: {score}")

数据混合策略

不同领域的数据需要按特定比例混合,才能训练出全面均衡的模型。数据混合策略直接影响模型的通用能力和专项能力。

典型数据混合比例

以下是 DeepSeek 类模型训练中常用的数据混合比例参考:

数据类别 建议比例 说明
通用对话 30-40% 日常对话、问答、闲聊,保证基础对话能力
代码生成 15-20% Python/JS/Java/C++ 等主流语言代码
数学推理 10-15% 代数、几何、微积分、概率统计
逻辑推理 10-15% 逻辑谜题、推理链、多步推理
创意写作 5-10% 诗歌、故事、文案、剧本
专业领域 5-10% 医学、法律、金融等垂直领域
安全对齐 3-5% 拒绝有害请求、价值观对齐

数据混合代码实现

import random from typing import List, Dict def mix_datasets(datasets: Dict[str, List], ratios: Dict[str, float], total_size: int = 10000, shuffle: bool = True): """按比例混合多个数据集""" # 验证比例总和 total_ratio = sum(ratios.values()) if abs(total_ratio - 1.0) > 0.01: print(f"警告: 比例总和为 {total_ratio:.2f},建议调整为 1.0") mixed = [] stats = {} for name, ratio in ratios.items(): if name not in datasets: continue # 按比例采样 sample_size = int(total_size * ratio) available = datasets[name] # 如果数据不足,全部使用 if len(available) < sample_size: sampled = available.copy() print(f" {name}: 数据不足,使用全部 {len(available)} 条") else: sampled = random.sample(available, sample_size) # 添加来源标记 for s in sampled: s["source_dataset"] = name mixed.extend(sampled) stats[name] = len(sampled) if shuffle: random.shuffle(mixed) print(f"\n数据集混合完成,共 {len(mixed)} 条") for name, count in stats.items(): print(f" {name}: {count} 条 ({count/len(mixed)*100:.1f}%)") return mixed # 使用示例 datasets = { "general_chat": [...], # 通用对话数据 "code": [...], # 代码数据 "math": [...], # 数学数据 "reasoning": [...], # 推理数据 "creative": [...], # 创意写作数据 } ratios = { "general_chat": 0.35, "code": 0.20, "math": 0.15, "reasoning": 0.15, "creative": 0.10, "safety": 0.05, } training_data = mix_datasets(datasets, ratios, total_size=10000)

数据退火策略

数据退火(Data Annealing)是一种训练策略,在训练过程中动态调整数据混合比例。早期使用多样化数据训练基础能力,后期逐步增加高质量、高难度数据的比例:

  • 预热阶段(前 20% 训练步数):通用对话数据占 50%,简单指令为主,帮助模型建立基础对话能力
  • 主训练阶段(20%-80% 训练步数):逐步增加代码和推理数据比例,引入中等难度任务
  • 退火阶段(最后 20% 训练步数):大幅增加高质量、高难度数据,减少通用对话数据,提升模型在复杂任务上的表现

DeepSeek 团队在训练中使用了类似的退火策略,这也是 DeepSeek-R1 在推理能力上表现突出的重要原因之一。

大规模数据处理

当数据量达到百万级甚至亿级时,单机处理已不可行。本章介绍使用 Spark、Ray 等分布式框架进行大规模数据处理,以及数据版本管理。

使用 Ray 并行处理

Ray 是一个轻量级的分布式计算框架,特别适合 Python 生态的数据处理任务:

import ray from typing import List, Dict # 初始化 Ray ray.init(num_cpus=8) # 将处理函数声明为 Ray 远程任务 @ray.remote def process_batch(batch: List[Dict]) -> List[Dict]: """处理一批数据:清洗、去重、质量过滤""" filter = DataQualityFilter() # 假设这里包含第 3 章中定义的所有过滤逻辑 return filter.apply_all(batch) def parallel_process(samples: List[Dict], batch_size=1000): """使用 Ray 并行处理大规模数据""" # 分批次 batches = [samples[i:i+batch_size] for i in range(0, len(samples), batch_size)] # 并行提交任务 futures = [process_batch.remote(batch) for batch in batches] # 收集结果 results = [] for future in futures: results.extend(ray.get(future)) print(f"并行处理完成: {len(samples)} -> {len(results)} 条") return results # 使用示例 # results = parallel_process(all_samples, batch_size=2000) # ray.shutdown()

使用 PySpark 分布式处理

PySpark 是处理 TB 级大规模数据的标准工具,支持 DataFrame API 和 SQL 操作:

from pyspark.sql import SparkSession from pyspark.sql.functions import col, length, udf from pyspark.sql.types import BooleanType # 创建 Spark 会话 spark = SparkSession.builder \ .appName("DeepSeekDataProcessing") \ .config("spark.sql.adaptive.enabled", "true") \ .getOrCreate() # 读取大规模 JSON 数据 df = spark.read.json("hdfs://data/deepseek/*.jsonl") # 基础过滤 df_clean = df.filter( (length(col("instruction")) > 20) & (length(col("output")) > 20) & (length(col("instruction")) < 8000) ) # 使用 SQL 进行复杂过滤 df_clean.createOrReplaceTempView("data") df_quality = spark.sql(""" SELECT *, LENGTH(instruction) + LENGTH(output) AS total_length FROM data WHERE instruction IS NOT NULL AND output IS NOT NULL AND LENGTH(output) / LENGTH(instruction) > 0.5 """) # 统计信息 print(f"处理前: {df.count()} 条") print(f"处理后: {df_quality.count()} 条") # 写入处理后的数据 df_quality.write.mode("overwrite").json("hdfs://data/deepseek_cleaned/")

流式处理大文件

对于超大的 JSONL 文件,使用流式处理避免内存溢出:

import json import ijson import gzip def stream_process_jsonl(filepath, output_path, batch_size=10000): """流式处理大型 JSONL 文件""" opener = gzip.open if filepath.endswith('.gz') else open batch = [] total_processed = 0 total_written = 0 filter = DataQualityFilter() with opener(filepath, 'r', encoding='utf-8') as infile, \ open(output_path, 'w', encoding='utf-8') as outfile: for line in infile: try: sample = json.loads(line.strip()) batch.append(sample) total_processed += 1 except json.JSONDecodeError: continue if len(batch) >= batch_size: # 批量处理 cleaned = filter.apply_all(batch) for s in cleaned: outfile.write(json.dumps(s, ensure_ascii=False) + '\n') total_written += len(cleaned) batch = [] if total_processed % 100000 == 0: print(f"已处理 {total_processed} 条, 写入 {total_written} 条") # 处理剩余批次 if batch: cleaned = filter.apply_all(batch) for s in cleaned: outfile.write(json.dumps(s, ensure_ascii=False) + '\n') total_written += len(cleaned) print(f"处理完成: {total_processed} -> {total_written} 条")

数据版本管理(DVC)

DVC(Data Version Control)是数据科学的 Git,用于管理数据集的版本:

# 初始化 DVC dvc init # 跟踪数据集 dvc add data/training_data.jsonl # 提交到 Git git add data/training_data.jsonl.dvc data/.gitignore git commit -m "添加初始训练数据集 v1.0" # 数据更新后,创建新版本 dvc add data/training_data.jsonl git add data/training_data.jsonl.dvc git commit -m "更新训练数据: 新增 10K 代码数据" # 切换到历史版本 git checkout "<commit-hash>" dvc checkout # 查看数据变更历史 dvc diff HEAD~1

生产级数据管线

将前面所有步骤整合为一条完整的生产级数据管线,包含持续数据采集、自动化质量监控、数据漂移检测和完整的 Pipeline 代码。

完整数据管线架构

  1. 数据采集层:定时任务抓取公开数据集更新、爬取新内容、调用合成数据 API
  2. 数据清洗层:质量过滤、去重、敏感信息过滤、格式标准化
  3. 数据增强层:Self-Instruct 生成、Evol-Instruct 进化、反向翻译
  4. 质量评估层:多样性评分、困难度评估、自动质量打分
  5. 数据混合层:按比例混合、数据退火策略、版本管理
  6. 监控告警层:数据漂移检测、质量趋势监控、异常告警

完整 Pipeline 代码

"""DeepSeek 数据工程管线 — 完整实现""" import json import os import logging from datetime import datetime from pathlib import Path from typing import List, Dict, Optional logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s') logger = logging.getLogger(__name__) class DataPipeline: """DeepSeek 数据工程管线""" def __init__(self, config: Dict): self.config = config self.data_dir = Path(config.get("data_dir", "./data")) self.output_dir = Path(config.get("output_dir", "./output")) self.data_dir.mkdir(exist_ok=True) self.output_dir.mkdir(exist_ok=True) self.metrics = {} def step1_collect(self) -> List[Dict]: """第1步:数据采集""" logger.info("Step 1: 数据采集") all_data = [] # 加载公开数据集 from datasets import load_dataset try: alpaca = load_dataset("tatsu-lab/alpaca", split="train") all_data.extend([dict(s) for s in alpaca]) logger.info(f" 加载 Alpaca: {len(alpaca)} 条") except Exception as e: logger.warning(f" 加载 Alpaca 失败: {e}") # 加载本地 JSONL 文件 for f in self.data_dir.glob("*.jsonl"): with open(f, 'r', encoding='utf-8') as fh: for line in fh: try: all_data.append(json.loads(line.strip())) except json.JSONDecodeError: continue logger.info(f" 加载 {f.name}") self.metrics["collected"] = len(all_data) return all_data def step2_clean(self, data: List[Dict]) -> List[Dict]: """第2步:数据清洗与去重""" logger.info("Step 2: 数据清洗与去重") # 长度过滤 def is_valid(s): inst = s.get("instruction", "") out = s.get("output", "") return 20 <= len(inst) <= 8000 and len(out) >= 20 cleaned = [s for s in data if is_valid(s)] logger.info(f" 长度过滤: {len(data)} -> {len(cleaned)} 条") # 去重(基于 instruction 的精确去重) seen = set() deduped = [] for s in cleaned: key = s.get("instruction", "").strip().lower() if key not in seen: seen.add(key) deduped.append(s) logger.info(f" 去重: {len(cleaned)} -> {len(deduped)} 条") self.metrics["cleaned"] = len(deduped) return deduped def step3_quality_check(self, data: List[Dict]) -> List[Dict]: """第3步:质量评估与过滤""" logger.info("Step 3: 质量评估") # 过滤 output 过短或包含拒绝模式的数据 reject_patterns = [ "抱歉,我无法", "作为一个AI", "对不起", "I cannot", "I'm sorry", "as an AI", ] def is_quality(s): output = s.get("output", "") if len(output) < 50: return False for p in reject_patterns: if p in output[:100]: return False return True quality = [s for s in data if is_quality(s)] logger.info(f" 质量过滤: {len(data)} -> {len(quality)} 条") self.metrics["quality_passed"] = len(quality) return quality def step4_export(self, data: List[Dict], filename: str = None): """第4步:导出处理后的数据""" logger.info("Step 4: 导出数据") if filename is None: timestamp = datetime.now().strftime("%Y%m%d_%H%M%S") filename = f"deepseek_data_{timestamp}.jsonl" output_path = self.output_dir / filename with open(output_path, 'w', encoding='utf-8') as f: for s in data: f.write(json.dumps(s, ensure_ascii=False) + '\n') logger.info(f" 导出到: {output_path} ({len(data)} 条)") return output_path def run(self): """运行完整管线""" logger.info("=" * 50) logger.info("DeepSeek 数据工程管线启动") logger.info("=" * 50) start = datetime.now() data = self.step1_collect() data = self.step2_clean(data) data = self.step3_quality_check(data) output_path = self.step4_export(data) elapsed = (datetime.now() - start).total_seconds() # 输出最终报告 logger.info("=" * 50) logger.info("管线运行完成") logger.info(f" 总耗时: {elapsed:.1f} 秒") for k, v in self.metrics.items(): logger.info(f" {k}: {v}") logger.info(f" 输出文件: {output_path}") logger.info("=" * 50) return data # ========== 主程序 ========== if __name__ == "__main__": config = { "data_dir": "./data", "output_dir": "./output", } pipeline = DataPipeline(config) pipeline.run()

数据漂移检测

数据漂移(Data Drift)是指数据分布随时间发生的变化。在生产环境中,需要持续监控数据质量是否偏离预期:

import numpy as np from scipy.stats import ks_2samp def detect_data_drift(reference_data: List[Dict], current_data: List[Dict], threshold=0.05): """检测当前数据是否相对于参考数据发生漂移""" # 1. 指令长度分布对比 ref_lengths = [len(s.get("instruction", "")) for s in reference_data] cur_lengths = [len(s.get("instruction", "")) for s in current_data] stat, p_value = ks_2samp(ref_lengths, cur_lengths) drift_detected = p_value < threshold print(f"KS 检验: stat={stat:.4f}, p={p_value:.4f}") print(f"平均长度: 参考={np.mean(ref_lengths):.0f}, 当前={np.mean(cur_lengths):.0f}") if drift_detected: print("警告: 检测到数据漂移!") else: print("数据分布正常,未检测到显著漂移。") return { "drift_detected": drift_detected, "p_value": p_value, "ref_mean": np.mean(ref_lengths), "cur_mean": np.mean(cur_lengths), }

自动化质量监控

使用定时任务持续监控数据质量指标:

import schedule import time def monitoring_job(): """定时监控任务""" logger.info("[监控] 开始数据质量检查...") # 运行管线 pipeline = DataPipeline(config) try: data = pipeline.run() # 检查质量指标 quality_ratio = pipeline.metrics.get("quality_passed", 0) / \ max(pipeline.metrics.get("collected", 1), 1) if quality_ratio < 0.5: logger.warning(f"[告警] 数据质量通过率过低: {quality_ratio:.2%}") else: logger.info(f"[正常] 数据质量通过率: {quality_ratio:.2%}") except Exception as e: logger.error(f"[错误] 监控任务失败: {e}") # 配置定时任务:每天凌晨 2 点执行 # schedule.every().day.at("02:00").do(monitoring_job) # # while True: # schedule.run_pending() # time.sleep(60)

运行 Pipeline

# 1. 准备数据目录 mkdir data # 将你的原始数据文件放入 data/ 目录 # 2. 安装依赖 pip install datasets schedule # 3. 运行管线 python data_pipeline.py # 4. 查看输出 # 处理后的数据在 output/ 目录 ls output/

生产环境建议

  • 使用 Apache Airflow 或 Prefect 编排管线任务
  • 将质量指标推送到 Prometheus + Grafana 进行可视化监控
  • 配置告警规则:数据质量通过率低于阈值、数据量异常波动等
  • 使用 DVC 或 LakeFS 进行数据版本管理,确保可复现
  • 定期人工抽检数据质量,与自动评分交叉验证

DeepSeek 数据工程常见问题

数据质量更重要还是数据数量更重要? +
在计算资源有限的情况下,数据质量远重要于数据数量。高质量的 10K 数据训练的模型,往往优于低质量的 100K 数据。DeepSeek 团队的经验表明,精心筛选和标注的数据集能让模型在指令遵循和推理能力上有质的飞跃。建议先保证质量,再考虑扩充数量。
SFT 数据最少需要多少条? +
对于高质量的 SFT 数据,10K-50K 条即可让模型展现出不错的指令遵循能力。如果追求全面能力,建议 100K-500K 条。关键是数据要覆盖多样化的任务类型和难度级别,而非单纯追求数量。Alpaca 的 52K 数据就是一个很好的起点。
如何判断数据是否被 DeepSeek 训练过? +
可以通过构造模型训练数据中不太可能出现的特异性问题来测试。例如,使用特定日期后的事件、特定格式的虚构数据等。如果模型能准确回答,说明可能有数据泄露。但这种方法不是绝对可靠的,最准确的方式是参考官方技术报告中的训练数据说明。
合成数据生成有什么风险? +
合成数据的主要风险包括:1) 模型幻觉被放大——如果生成模型本身有错误,合成数据会放大这些错误;2) 多样性不足——合成数据往往缺乏真实数据的多样性;3) 模式坍缩——多次迭代合成可能导致数据分布越来越窄。建议合成数据与真实数据混合使用,并进行严格的质量过滤。
DeepSeek-R1 的 thinking 数据如何准备? +
DeepSeek-R1 的 thinking 数据需要包含显式的推理过程,用 <think> 和 </think> 标签包裹。可以使用强模型(如 DeepSeek-V3)为数学和推理问题生成 CoT 推理链,然后人工审核质量。关键是要确保推理过程正确、完整、可验证。不建议使用 R1 自己生成的推理数据再训练 R1,可能导致模式坍缩。
数据清洗会丢失有价值的信息吗? +
过于激进的数据清洗确实可能丢失有价值的信息。建议采用渐进式清洗策略:先用宽松的规则过滤明显低质量的数据,然后在训练过程中根据模型表现动态调整过滤规则。对于边界情况,可以保留样本并标记为低置信度,而非直接丢弃。定期对丢弃的数据进行人工抽检,评估清洗规则的合理性。

DeepSeek 相关教程

深入学习 DeepSeek 模型的使用、部署和生态工具。

每日精选 Skill 推荐,免费送到你邮箱

输入邮箱,每天接收一个精选 AI Agent 技能推荐。完全免费,持续更新。

完全免费,取消任意时间。我们不会发送垃圾邮件。