【大模型】- 构建完整的 LLM 管道

算法

构建完整的 LLM 管道

整合所有组件,构建生产级系统

类型: 学习 | 语言: Python | 🏷 前置:《完整 LLM 管道》(本系列第 12 篇)

学习目标

  • 整合预训练、微调和推理
  • 构建完整的训练管道
  • 实现模型服务
  • 监控和维护生产系统
  • 处理常见的生产问题

管道架构

1
2
数据收集 → 清洗 → Tokenization → 预训练 → 评估 → 
微调(SFT/RLHF)→ 量化 → 推理服务 → 监控

阶段 1:数据管道

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
class DataPipeline:
def __init__(self, config):
self.config = config
self.tokenizer = AutoTokenizer.from_pretrained(config["model_name"])

def collect_data(self):
"""收集数据"""
# 从多个来源加载
datasets = []

# 网络数据
if "web" in self.config["sources"]:
web_data = load_dataset("c4", split="train")
datasets.append(web_data)

# 书籍数据
if "books" in self.config["sources"]:
book_data = load_dataset("bookcorpus", split="train")
datasets.append(book_data)

# 合并数据集
combined = concatenate_datasets(datasets)
return combined

def preprocess(self, dataset):
"""预处理数据"""
def clean_text(example):
# 清理文本
text = example["text"]
text = text.strip()
text = re.sub(r'\s+', ' ', text)
return {"text": text}

# 清理
dataset = dataset.map(clean_text)

# 分词
def tokenize(example):
return self.tokenizer(
example["text"],
truncation=True,
max_length=self.config["max_length"]
)

dataset = dataset.map(tokenize, batched=True)

return dataset

def create_dataloader(self, dataset):
"""创建 DataLoader"""
dataset.set_format(type="torch", columns=["input_ids", "attention_mask"])

return DataLoader(
dataset,
batch_size=self.config["batch_size"],
shuffle=True,
num_workers=4
)

阶段 2:训练管道

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
class TrainingPipeline:
def __init__(self, config):
self.config = config
self.model = self.load_model()
self.optimizer = self.create_optimizer()

def load_model(self):
"""加载模型"""
model = GPT(self.config)

# 检查点恢复
if self.config.get("checkpoint"):
checkpoint = torch.load(self.config["checkpoint"])
model.load_state_dict(checkpoint["model"])

return model

def create_optimizer(self):
"""创建优化器"""
return torch.optim.AdamW(
self.model.parameters(),
lr=self.config["lr"],
weight_decay=self.config["weight_decay"]
)

def train_epoch(self, dataloader):
"""训练一个 epoch"""
self.model.train()
total_loss = 0

for batch in dataloader:
input_ids = batch["input_ids"].to(self.config["device"])
attention_mask = batch["attention_mask"].to(self.config["device"])
labels = input_ids.clone()

# 前向传播
logits, loss = self.model(input_ids, labels=labels)

# 反向传播
loss.backward()
torch.nn.utils.clip_grad_norm_(
self.model.parameters(),
self.config["max_grad_norm"]
)

self.optimizer.step()
self.optimizer.zero_grad()

total_loss += loss.item()

return total_loss / len(dataloader)

def train(self, train_dataloader, val_dataloader):
"""完整训练流程"""
best_val_loss = float("inf")

for epoch in range(self.config["epochs"]):
# 训练
train_loss = self.train_epoch(train_dataloader)

# 验证
val_loss = self.evaluate(val_dataloader)

print(f"Epoch {epoch + 1}:")
print(f" Train Loss: {train_loss:.4f}")
print(f" Val Loss: {val_loss:.4f}")

# 保存检查点
if val_loss < best_val_loss:
best_val_loss = val_loss
self.save_checkpoint(f"best_model.pt")

# 定期保存
if (epoch + 1) % self.config["save_interval"] == 0:
self.save_checkpoint(f"checkpoint-{epoch + 1}.pt")

def save_checkpoint(self, path):
"""保存检查点"""
torch.save({
"model": self.model.state_dict(),
"optimizer": self.optimizer.state_dict(),
"config": self.config,
}, path)

阶段 3:评估管道

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
class EvaluationPipeline:
def __init__(self, model, tokenizer, config):
self.model = model
self.tokenizer = tokenizer
self.config = config

def evaluate_benchmarks(self, benchmarks):
"""评估多个基准"""
results = {}

for benchmark in benchmarks:
print(f"评估 {benchmark.name}...")

score = benchmark.evaluate(self.model, self.tokenizer)
results[benchmark.name] = score

return results

def evaluate_safety(self, test_cases):
"""评估安全性"""
safety_scores = []

for test_case in test_cases:
response = self.generate(test_case["prompt"])

# 检查是否有害内容
is_safe = self.check_safety(response)
safety_scores.append(is_safe)

return sum(safety_scores) / len(safety_scores)

def generate(self, prompt, **kwargs):
"""生成文本"""
inputs = self.tokenizer(prompt, return_tensors="pt").to(self.model.device)

with torch.no_grad():
outputs = self.model.generate(
**inputs,
**kwargs
)

return self.tokenizer.decode(outputs[0], skip_special_tokens=True)

def check_safety(self, response):
"""检查安全性"""
# 简单的安全检查
unsafe_patterns = [
"暴力", "仇恨", "歧视", "危险",
"illegal", "harmful", "dangerous"
]

for pattern in unsafe_patterns:
if pattern in response.lower():
return False

return True

阶段 4:推理服务

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
class LLMService:
def __init__(self, model_path, config):
self.config = config
self.model = self.load_model(model_path)
self.tokenizer = AutoTokenizer.from_pretrained(config["model_name"])

# 优化
if config.get("use_flash_attention"):
self.model = self.optimize_model()

# 监控
self.monitor = InferenceMonitor()

def load_model(self, model_path):
"""加载模型"""
model = AutoModelForCausalLM.from_pretrained(
model_path,
torch_dtype=torch.float16,
device_map="auto"
)

if self.config.get("quantized"):
model = self.quantize_model(model)

return model

def optimize_model(self):
"""优化模型"""
# 使用 Flash Attention
self.model = self.model.to(dtype=torch.bfloat16)

# 使用 torch.compile
if self.config.get("use_torch_compile"):
self.model = torch.compile(self.model)

return self.model

def quantize_model(self, model):
"""量化模型"""
from transformers import BitsAndBytesConfig

bnb_config = BitsAndBytesConfig(
load_in_4bit=True,
bnb_4bit_use_double_quant=True,
bnb_4bit_quant_type="nf4",
)

return AutoModelForCausalLM.from_pretrained(
self.config["model_name"],
quantization_config=bnb_config,
device_map="auto"
)

async def generate(self, request):
"""生成文本"""
start_time = time.time()

inputs = self.tokenizer(
request.prompt,
return_tensors="pt"
).to(self.model.device)

with torch.no_grad():
outputs = self.model.generate(
**inputs,
max_new_tokens=request.max_new_tokens,
temperature=request.temperature,
top_p=request.top_p,
do_sample=True
)

response = self.tokenizer.decode(
outputs[0],
skip_special_tokens=True
)

tokens_generated = outputs.shape[1] - inputs["input_ids"].shape[1]
time_taken = time.time() - start_time

self.monitor.record_request(tokens_generated, time_taken)

return {
"response": response,
"tokens_generated": tokens_generated,
"time_taken": time_taken
}

监控和维护

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
class LLMMonitor:
def __init__(self, config):
self.config = config
self.metrics = {}

# 初始化监控
if config.get("use_wandb"):
import wandb
wandb.init(project="llm-production")

def log_metrics(self, metrics):
"""记录指标"""
self.metrics.update(metrics)

if self.config.get("use_wandb"):
import wandb
wandb.log(metrics)

def check_health(self, service):
"""健康检查"""
health_status = {
"model_loaded": service.model is not None,
"gpu_available": torch.cuda.is_available(),
"memory_usage": self.get_memory_usage(),
"request_rate": self.get_request_rate(),
}

return health_status

def get_memory_usage(self):
"""获取内存使用情况"""
if torch.cuda.is_available():
return {
"allocated": torch.cuda.memory_allocated(),
"cached": torch.cuda.memory_reserved(),
"max_allocated": torch.cuda.max_memory_allocated()
}
return {}

def get_request_rate(self):
"""获取请求速率"""
# 计算每秒请求数
return self.metrics.get("requests_per_second", 0)

部署配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
# config.yaml
training:
model_name: "gpt2"
batch_size: 32
learning_rate: 3e-4
epochs: 10
save_interval: 2

evaluation:
benchmarks:
- name: "mmlu"
num_samples: 1000
- name: "hellaswag"
num_samples: 1000

inference:
model_path: "./checkpoints/best_model.pt"
quantized: true
use_flash_attention: true
max_batch_size: 32
max_sequence_length: 2048

monitoring:
use_wandb: true
log_interval: 100
health_check_interval: 60

完整管道

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
def main():
# 加载配置
config = load_config("config.yaml")

# 数据管道
data_pipeline = DataPipeline(config["data"])
dataset = data_pipeline.collect_data()
dataset = data_pipeline.preprocess(dataset)
dataloader = data_pipeline.create_dataloader(dataset)

# 训练管道
training_pipeline = TrainingPipeline(config["training"])
training_pipeline.train(dataloader, dataloader) # 简化:使用相同的数据

# 评估管道
eval_pipeline = EvaluationPipeline(
training_pipeline.model,
data_pipeline.tokenizer,
config["evaluation"]
)
results = eval_pipeline.evaluate_benchmarks([
MMLUBenchmark(),
HellaSwagBenchmark()
])

# 推理服务
service = LLMService(
"./checkpoints/best_model.pt",
config["inference"]
)

# 启动服务
app = FastAPI()

@app.post("/generate")
async def generate(request: GenerationRequest):
return await service.generate(request)

# 监控
monitor = LLMMonitor(config["monitoring"])

uvicorn.run(app, host="0.0.0.0", port=8000)

if __name__ == "__main__":
main()

总结

完整的 LLM 管道整合了数据处理、训练、评估、推理和监控。每个阶段都需要仔细设计和优化。生产系统需要处理错误、监控性能并支持扩展。

下一步

下一课将介绍开源模型架构,分析主流 LLM 的设计。

📚 本文改编自 AI Engineering from Scratch(MIT License · 作者 Rohit Ghumare),中文内容来自官方中文镜像。原课程共 503 课 · 20 阶段 · 免费开源,教程网站见 aiengineeringfromscratch.com

  • 标题: 【大模型】- 构建完整的 LLM 管道
  • 作者:
  • 创建于 : 2026-08-19 09:13:00
  • 更新于 : 2026-08-21 16:20:11
  • 链接: https://sxl-space.tk/2026/08/19/010_LLM/010_LLM-13-CompleteLLMPipeline/
  • 版权声明: 版权所有 © 宋,禁止转载。