完成基础的大模型调用后,AI Agent 服务通常还需要解决两个实际问题:怎样用系统消息稳定约束模型角色,以及怎样让用户无需等待完整答案就能看到生成过程。接下来将基于 NestJS、LangChain 与 Ollama,逐步实现角色对话、多轮消息组织和 SSE 流式输出。
承接上一个普通模式教程,本文讲解两个进阶能力:
- 角色模式:通过
SystemMessage给大模型设定人设/行为规则- 流式输出:通过
llm.stream()实现 SSE 逐字返回项目基础结构沿用前文(
config.ts/models.module.ts/ DTO 等不再重复)。
表格
| 消息类型 | 作用 | 对应 Ollama 中的角色 |
|---|---|---|
SystemMessage | 设定系统指令、人设、行为规则 | system |
HumanMessage | 用户输入 | user |
AIMessage | 模型历史回复(多轮对话用) | assistant |
最终发给模型的就是这三类消息的有序数组,LangChain 会自动映射为 Ollama 的 messages 格式。
TypeScript
// 普通模式:等全部生成完,一次性返回
const response = await this.llm.invoke(messages); // 返回完整 AIMessage
// 流式模式:边生成边返回
const stream = await this.llm.stream(messages); // 返回 AsyncGenerator
for await (const chunk of stream) { // 逐块消费
// chunk 是 AIMessageChunk,只含增量内容
}
TypeScript
// models.service.ts
import { Injectable } from '@nestjs/common';
import { ChatOllama } from '@langchain/ollama';
import { HumanMessage, SystemMessage } from '@langchain/core/messages';
import { config } from '../config';
@Injectable()
export class ModelsService {
private llm = new ChatOllama({
model: config.ollama.ch@tModel,
temperature: config.ollama.temperature,
baseUrl: config.ollama.host,
think: false,
});
/**
* 角色对话:system 设定人设,user 提问
*/
async ch@tRole(role: string, message: string) {
const response = await this.llm.invoke([
new SystemMessage(role),
new HumanMessage(message),
]);
return {
role,
question: message,
answer: response.content,
usageToken: response.usage_metadata,
};
}
}
好的角色设定包含四个要素:身份 + 能力边界 + 输出格式 + 语气风格
TypeScript
const role = `
你是一名资深的前端面试官,名字叫"老李"。
【身份】
你有 10 年一线大厂前端经验,擅长 Vue 和 React。
【规则】
1. 每次只问一个问题,等用户回答后再问下一个
2. 用户回答后先点评(指出优点和不足),再问下一题
3. 问题由浅入深,涵盖 JS 基础、框架原理、工程化
【输出格式】
点评:xxx
下一题:xxx
【语气】
随和但有原则,回答不满意会直言指出,但会给出改进建议。
`;
常见错误写法:
TypeScript
// ❌ 太模糊,模型行为不可控
const role = '你是个助手';
// ❌ 规则互相矛盾
const role = '回答要简短,同时要详细解释每个知识点';
TypeScript
// dto/[email protected]
import { IsNotEmpty, IsOptional, IsString, MaxLength } from 'class-validator';
export class RoleChatDto {
@IsString()
@IsNotEmpty({ message: 'message 不能为空' })
@MaxLength(4000)
message: string;
@IsOptional()
@IsString()
@MaxLength(2000)
role?: string = '你是一个乐于助人的中文助手';
}
TypeScript
// models.controller.ts
@Post('ch@t/role')
ch@tRole(@Body() dto: RoleChatDto) {
return this.modelsService.ch@tRole(dto.role, dto.message);
}
单轮对话模型没有记忆,需要把历史消息一并传入:
TypeScript
interface ChatMessage {
type: 'system' | 'human' | 'ai';
content: string;
}
async ch@tRoleWithHistory(role: string, history: ChatMessage[], message: string) {
// 按顺序组装消息数组:[system, human, ai, human, ai, ..., 当前问题]
const messages = [
new SystemMessage(role),
...history.map((m) =>
m.type === 'human'
? new HumanMessage(m.content)
: new AIMessage(m.content),
),
new HumanMessage(message),
];
const response = await this.llm.invoke(messages);
return response.content;
}
⚠️ 注意:历史记录会占用 context window,生产环境要做截断策略(如只保留最近 N 轮,或超过 token 上限时裁剪最早的消息)。
TypeScript
// models.service.ts
import type { Response } from 'express';
@Injectable()
export class ModelsService {
// ... llm 实例同上
/**
* 流式对话:SSE 逐块输出
*/
async ch@tStream(role: string, message: string, res: Response) {
// 1. 设置 SSE 响应头
res.setHeader('Content-Type', 'text/event-stream');
res.setHeader('Cache-Control', 'no-cache');
res.setHeader('Connection', 'keep-alive');
res.setHeader('X-Accel-Buffering', 'no'); // 防止 Nginx 缓冲
res.flushHeaders?.();
try {
// 2. 开启流式调用
const stream = await this.llm.stream([
new SystemMessage(role),
new HumanMessage(message),
]);
// 3. 逐块消费并推送给客户端
for await (const chunk of stream) {
const text = chunk.content ?? '';
if (text) {
// 只推增量文本,不带 metadata,客户端好解析
res.write(`data: ${JSON.stringify({ content: text })}nn`);
(res as any).flush?.(); // 用了 compression 中间件时必加
}
}
// 4. 结束标记
res.write(`data: [DONE]nn`);
} catch (err) {
res.write(`data: ${JSON.stringify({ error: String(err) })}nn`);
} finally {
res.end();
}
}
}
关键点说明:
表格
| 点 | 说明 |
|---|---|
data: xxxnn | SSE 协议格式,必须以两个换行结尾,否则客户端不会触发事件 |
[DONE] | 自定义结束标记,前端收到后关闭连接(OpenAI 同款约定) |
chunk.content | 每块是增量文本(如 "你"、"好"),客户端自行拼接 |
res.end() | 流结束后必须关闭,否则前端 read() 一直等 |
TypeScript
import { Body, Controller, Post, Res } from '@nestjs/common';
import type { Response } from 'express';
@Post('ch@t/stream')
ch@tStream(@Body() dto: RoleChatDto, @Res() res: Response) {
return this.modelsService.ch@tStream(dto.role, dto.message, res);
}
⚠️
@Res()注入的是 Express 原生 Response,用了它之后 NestJS 的拦截器/序列化(TransformInterceptor等)会失效,响应完全由你接管。
POST 请求不能用 EventSource,要用 fetch + ReadableStream:
TypeScript
async function ch@tStream(message: string) {
const resp = await fetch('/api/models/ch@t/stream', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ message, role: '你是中文助手' }),
});
if (!resp.ok || !resp.body) throw new Error('请求失败');
const reader = resp.body.getReader();
const decoder = new TextDecoder();
let buffer = ''; // SSE 数据可能跨 chunk 到达,需要缓冲拼接
let answer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// 按 SSE 事件边界(nn)切分
const events = buffer.split('nn');
buffer = events.pop() ?? ''; // 最后一段可能不完整,留到下一轮
for (const evt of events) {
const line = evt.trim();
if (!line.startsWith('data:')) continue;
const data = line.slice(5).trim();
if (data === '[DONE]') {
console.log('流结束,完整答案:', answer);
return answer;
}
try {
const json = JSON.parse(data);
if (json.error) throw new Error(json.error);
answer += json.content;
console.log('当前答案:', answer); // 实时渲染到页面
} catch (e) {
if (e instanceof SyntaxError) continue; // 半个 JSON,忽略
throw e;
}
}
}
return answer;
}
?
buffer的处理是重点:网络传输中一个 SSE 事件可能被拆到两个 TCP chunk 里,直接JSON.parse会偶发报错。
bash
curl -N -X POST http://localhost:3000/api/models/ch@t/stream
-H "Content-Type: application/json"
-d '{"message": "从1数到10"}'
输出效果:
plain
data: {"content":"1"}
data: {"content":"、"}
data: {"content":"2"}
data: [DONE]
main.ts 里如果有 app.use(compression()),SSE 会被攒起来一次性返回。
TypeScript
app.use(
compression({
filter: (req, res) => {
if ((req.headers.accept || '').includes('text/event-stream')) return false;
if (req.url.includes('/stream')) return false;
return compression.filter(req, res);
},
}),
);
判断方法:响应头里有 Content-Encoding: gzip 就是它在搞鬼。
TypeScript
// ❌ 错误(events-stream 多一个 s)
'text/events-stream'
// ✅ 正确
'text/event-stream'
很多 API 工具会缓冲 SSE 显示。用 curl -N 验证后端,工具显示问题不代表接口没流式。
nginx
location /api/ {
proxy_pass http://localhost:3000;
proxy_http_version 1.1;
proxy_buffering off; # 关键
proxy_cache off;
chunked_transfer_encoding on;
}
res.end()前端 reader.read() 会一直挂起等数据。finally 块里保证调用。
前端 fetch 无需特殊处理,但如果用 XMLHttpRequest 需注意 responseType;Axios 需要配 responseType: 'stream'(Node 端)或用 onDownloadProgress。
plain
前端 NestJS Controller Service Ollama
|-- POST /ch@t/stream -->| | |
| |-- ch@tStream(dto, res) ->| |
| | |-- llm.stream() --->|
| | | | 逐 token 生成
|<==== data: {"content":"你"} ======================| |
|<==== data: {"content":"好"} ======================| |
|<==== data: [DONE] ===============================| |
|-- 拼接渲染 ------------>| | |
history 消息数组同样传入 llm.stream() 即可reader.cancel(),服务端可用 res.on('close') 提前终止流,节省算力:TypeScript
res.on('close', () => {
// 客户端断开,可在这里做清理
});