第 17 章 · 流式数据、元数据与消息持久化
本章目标:
- 掌握
createUIMessageStream+writer.write向前端流式发送自定义数据(data parts、sources、transient parts)- 理解 data part 的 reconciliation(同 ID 更新)机制与
onData回调处理- 学会用 message metadata 在消息级别附加时间戳、模型信息、token 用量等数据
- 实现聊天消息的持久化:存储/加载/校验 UIMessages,以及服务端消息 ID 生成
17.1 为什么需要流式自定义数据
在模型的响应之外,经常还需要向前端发送额外数据。例如:发送状态信息、在存库后回传 message id、或者引用语言模型正在提及的内容。
AI SDK 提供了几个辅助函数让你把额外数据流式发给客户端并挂到 UIMessage 的 parts 数组上:
createUIMessageStream:创建一个 data stream;createUIMessageStreamResponse:创建流式返回数据的响应对象;pipeUIMessageStreamToResponse:把 data stream 管道输出到 server response 对象。
数据以 Server-Sent Events 形式随响应流传送。
类型安全的数据流定义
先定义带 data part schema 的自定义消息类型:
import { UIMessage } from 'ai';
// Define your custom message type with data part schemas
export type MyUIMessage = UIMessage<
never, // metadata type
{
weather: {
city: string;
weather?: string;
status: 'loading' | 'success';
};
notification: {
message: string;
level: 'info' | 'warning' | 'error';
};
} // data parts type
>;17.2 从服务端流式发送数据
在服务端路由中创建 UIMessageStream 并交给 createUIMessageStreamResponse 返回:
import {
convertToModelMessages,
createUIMessageStream,
createUIMessageStreamResponse,
streamText,
toUIMessageStream,
} from 'ai';
import { createGateway } from 'ai';
import type { MyUIMessage } from '@/ai/types';
const gateway = createGateway({ apiKey: process.env.AI_GATEWAY_API_KEY ?? '' });
export async function POST(req: Request) {
const { messages } = await req.json();
const stream = createUIMessageStream<MyUIMessage>({
execute: ({ writer }) => {
// 1. Send initial status (transient - won't be added to message history)
writer.write({
type: 'data-notification',
data: { message: 'Processing your request...', level: 'info' },
transient: true, // This part won't be added to message history
});
// 2. Send sources (useful for RAG use cases)
writer.write({
type: 'source',
value: {
type: 'source',
sourceType: 'url',
id: 'source-1',
url: 'https://weather.com',
title: 'Weather Data Source',
},
});
// 3. Send data parts with loading state
writer.write({
type: 'data-weather',
id: 'weather-1',
data: { city: 'San Francisco', status: 'loading' },
});
const result = streamText({
model: gateway('openai/gpt-5'),
messages: await convertToModelMessages(messages),
onEnd() {
// 4. Update the same data part (reconciliation)
writer.write({
type: 'data-weather',
id: 'weather-1', // Same ID = update existing part
data: {
city: 'San Francisco',
weather: 'sunny',
status: 'success',
},
});
// 5. Send completion notification (transient)
writer.write({
type: 'data-notification',
data: { message: 'Request completed', level: 'info' },
transient: true, // Won't be added to message history
});
},
});
writer.merge(toUIMessageStream({ stream: result.stream }));
},
});
return createUIMessageStreamResponse({ stream });
}💡 自定义后端(如 Python/FastAPI)同样可以按 UI Message Stream 协议发送流式数据。
三类可流式传输的数据
持久化 data parts——会进入消息历史,出现在 message.parts 中:
writer.write({
type: 'data-weather',
id: 'weather-1', // Optional: enables reconciliation
data: { city: 'San Francisco', status: 'loading' },
});sources——适合 RAG 场景展示引用的文档或链接(见上方代码第 2 步)。
瞬态(transient)data parts——发送到客户端但不进入消息历史,只能通过 useChat 的 onData 回调访问:
// server
writer.write({
type: 'data-notification',
data: { message: 'Processing...', level: 'info' },
transient: true, // Won't be added to message history
});
// client
const [notification, setNotification] = useState();
const { messages } = useChat({
onData: ({ data, type }) => {
if (type === 'data-notification') {
setNotification({ message: data.message, level: data.level });
}
},
});Reconciliation:同 ID 数据部分自动更新
向同一个 id 写入时,客户端会自动 reconcile 并更新该部分。这能实现强大的动态体验:
- 协作式 artifacts——实时更新代码、文档或设计;
- 渐进式加载——loading 态平滑过渡为最终结果;
- 实时状态——更新进度条、计数器、状态指示器;
- 交互组件——随用户交互演化的 UI 元素。
reconciliation 自动发生:写入流时使用相同 id 即可。
17.3 客户端处理流式数据
onData 回调
onData 是处理流式数据的关键回调,尤其对瞬态 parts:
import { useChat } from '@ai-sdk/react';
import type { MyUIMessage } from '@/ai/types';
const { messages } = useChat<MyUIMessage>({
api: '/api/chat',
onData: dataPart => {
// Handle all data parts as they arrive (including transient parts)
console.log('Received data part:', dataPart);
// Handle different data part types
if (dataPart.type === 'data-weather') {
console.log('Weather update:', dataPart.data);
}
// Handle transient notifications (ONLY available here, not in message.parts)
if (dataPart.type === 'data-notification') {
showToast(dataPart.data.message, dataPart.data.level);
}
},
});⚠️ 重要:瞬态 data parts 只能通过
onData回调获取,不会出现在message.parts里(因为不进入消息历史)。
渲染持久化数据部分
从 message parts 数组中过滤渲染 data parts:
const result = (
<>
{messages?.map(message => (
<div key={message.id}>
{/* Render weather data parts */}
{message.parts
.filter(part => part.type === 'data-weather')
.map((part, index) => (
<div key={index} className="weather-widget">
{part.data.status === 'loading' ? (
<>Getting weather for {part.data.city}...</>
) : (
<>
Weather in {part.data.city}: {part.data.weather}
</>
)}
</div>
))}
{/* Render sources */}
{message.parts
.filter(part => part.type === 'source')
.map((part, index) => (
<div key={index} className="source">
Source: <a href={part.url}>{part.title}</a>
</div>
))}
</div>
))}
</>
);典型应用场景:RAG 应用(流式发送来源文档)、实时状态(加载进度)、协作工具、用量分析、临时通知。
17.4 Message Metadata:消息级元数据
message metadata 允许你在消息级别附加自定义信息,适合追踪时间戳、模型信息、token 用量、用户上下文等「关于整条消息」的数据——这与构成消息内容的 data parts 形成互补。
定义元数据类型
先为类型安全定义 metadata schema:
import { UIMessage } from 'ai';
import { z } from 'zod';
// Define your metadata schema
export const messageMetadataSchema = z.object({
createdAt: z.number().optional(),
model: z.string().optional(),
totalTokens: z.number().optional(),
});
export type MessageMetadata = z.infer<typeof messageMetadataSchema>;
// Create a typed UIMessage
export type MyUIMessage = UIMessage<MessageMetadata>;服务端发送元数据
用 toUIMessageStream 的 messageMetadata 回调在不同流式阶段发送元数据:
import {
convertToModelMessages,
createUIMessageStreamResponse,
streamText,
toUIMessageStream,
} from 'ai';
import { createGateway } from 'ai';
import type { MyUIMessage } from '@/types';
const gateway = createGateway({ apiKey: process.env.AI_GATEWAY_API_KEY ?? '' });
export async function POST(req: Request) {
const { messages }: { messages: MyUIMessage[] } = await req.json();
const result = streamText({
model: gateway('openai/gpt-5'),
messages: await convertToModelMessages(messages),
});
return createUIMessageStreamResponse({
stream: toUIMessageStream({
stream: result.stream,
originalMessages: messages, // pass this in for type-safe return objects
messageMetadata: ({ part }) => {
// Send metadata when streaming starts
if (part.type === 'start') {
return {
createdAt: Date.now(),
model: 'your-model-id',
};
}
// Send additional metadata when streaming completes
if (part.type === 'finish') {
return {
totalTokens: part.totalUsage.totalTokens,
};
}
},
}),
});
}💡 要在
messageMetadata中获得类型安全的返回对象,请传入类型化为你 UIMessage 类型的originalMessages参数。
客户端访问元数据
通过 message.metadata 属性读取:
'use client';
import { useChat } from '@ai-sdk/react';
import { DefaultChatTransport } from 'ai';
import type { MyUIMessage } from '@/types';
export default function Chat() {
const { messages } = useChat<MyUIMessage>({
transport: new DefaultChatTransport({
api: '/api/chat',
}),
});
return (
<div>
{messages.map(message => (
<div key={message.id}>
<div>
{message.role === 'user' ? 'User: ' : 'AI: '}
{message.metadata?.createdAt && (
<span className="text-sm text-gray-500">
{new Date(message.metadata.createdAt).toLocaleTimeString()}
</span>
)}
</div>
{/* Render message content */}
{message.parts.map((part, index) =>
part.type === 'text' ? <div key={index}>{part.text}</div> : null,
)}
{/* Display additional metadata */}
{message.metadata?.totalTokens && (
<div className="text-xs text-gray-400">
{message.metadata.totalTokens} tokens
</div>
)}
</div>
))}
</div>
);
}metadata 的常见用途:时间戳、模型信息、token 用量、用户上下文、性能指标(首 token 时间)、质量指标(finish reason)。若需要生成过程中动态变化的数据,应改用 data parts。
17.5 消息持久化
能够存储和加载聊天消息对大多数 AI 聊天应用至关重要。本节演示如何结合 useChat 与 streamText 实现消息持久化。
创建新聊天
用户不带聊天 ID 访问聊天页时,需要先创建新聊天再重定向:
import { redirect } from 'next/navigation';
import { createChat } from '@util/chat-store';
export default async function Page() {
const id = await createChat(); // create a new chat
redirect(`/chat/${id}`); // redirect to chat page, see below
}示例 chat store 用文件系统存储聊天消息;真实应用应改用数据库或云存储,函数接口设计上可以无缝替换:
import { generateId } from 'ai';
import { existsSync, mkdirSync } from 'fs';
import { writeFile } from 'fs/promises';
import path from 'path';
// Treat chat IDs as opaque tokens before using them in file paths.
const chatIdRegex = /^[A-Za-z0-9_-]+$/;
export async function createChat(): Promise<string> {
const id = generateId(); // generate a unique chat ID
await writeFile(getChatFile(id), '[]'); // create an empty chat file
return id;
}
function getChatFile(id: string): string {
if (!chatIdRegex.test(id)) {
throw new Error('Invalid chat ID');
}
const chatDir = path.resolve(process.cwd(), '.chats');
const chatFile = path.resolve(chatDir, `${id}.json`);
// Defense in depth: keep the resolved file inside the chat directory.
if (!chatFile.startsWith(`${chatDir}${path.sep}`)) {
throw new Error('Invalid chat ID');
}
if (!existsSync(chatDir)) mkdirSync(chatDir, { recursive: true });
return chatFile;
}⚠️ 聊天 ID 可能来自 URL 或请求体,把它当作用作文件路径前的 opaque token 校验;resolved path 检查确保文件不会逃出
.chats目录。
加载已有聊天
import { UIMessage } from 'ai';
import { readFile } from 'fs/promises';
export async function loadChat(id: string): Promise<UIMessage[]> {
return JSON.parse(await readFile(getChatFile(id), 'utf8'));
}
// ... rest of the file服务端校验消息
当处理包含 tool calls、自定义 metadata 或 data parts 的历史消息时,应在发送给模型之前用 validateUIMessages 校验:
import {
convertToModelMessages,
createUIMessageStreamResponse,
streamText,
toUIMessageStream,
UIMessage,
validateUIMessages,
tool,
} from 'ai';
import { z } from 'zod';
import { loadChat, saveChat } from '@util/chat-store';
import { dataPartsSchema, metadataSchema } from '@util/schemas';
import { createGateway } from 'ai';
const gateway = createGateway({ apiKey: process.env.AI_GATEWAY_API_KEY ?? '' });
// Define your tools
const tools = {
weather: tool({
description: 'Get weather information',
inputSchema: z.object({
location: z.string(),
units: z.enum(['celsius', 'fahrenheit']),
}),
execute: async ({ location, units }) => {
/* tool implementation */
},
}),
// other tools
};
export async function POST(req: Request) {
const { message, id } = await req.json();
// Load previous messages from database
const previousMessages = await loadChat(id);
// Append new message to previousMessages messages
const messages = [...previousMessages, message];
// Validate loaded messages against
// tools, data parts schema, and metadata schema
const validatedMessages = await validateUIMessages({
messages,
tools, // Ensures tool calls in messages match current schemas
dataPartsSchema,
metadataSchema,
});
const result = streamText({
model: gateway('openai/gpt-5-mini'),
messages: convertToModelMessages(validatedMessages),
tools,
});
return createUIMessageStreamResponse({
stream: toUIMessageStream({
stream: result.stream,
originalMessages: messages,
onEnd: ({ messages }) => {
saveChat({ chatId: id, messages });
},
}),
});
}数据库中的旧消息可能与当前 schema 不匹配,要优雅处理校验错误:
import {
convertToModelMessages,
streamText,
validateUIMessages,
TypeValidationError,
} from 'ai';
import { type MyUIMessage } from '@/types';
export async function POST(req: Request) {
const { message, id } = await req.json();
// Load and validate messages from database
let validatedMessages: MyUIMessage[];
try {
const previousMessages = await loadMessagesFromDB(id);
validatedMessages = await validateUIMessages({
// append the new message to the previous messages:
messages: [...previousMessages, message],
tools,
metadataSchema,
});
} catch (error) {
if (error instanceof TypeValidationError) {
// Log validation error for monitoring
console.error('Database messages validation failed:', error);
// Could implement message migration or filtering here
// For now, start with empty history
validatedMessages = [];
} else {
throw error;
}
}
// Continue with validated messages...
}展示聊天页面与保存消息
页面组件加载历史后传给聊天组件,后者用 id 和 initialMessages 初始化 useChat:
import { loadChat } from '@util/chat-store';
import Chat from '@ui/chat';
export default async function Page(props: { params: Promise<{ id: string }> }) {
const { id } = await props.params;
const messages = await loadChat(id);
return <Chat id={id} initialMessages={messages} />;
}'use client';
import { UIMessage, useChat } from '@ai-sdk/react';
import { DefaultChatTransport } from 'ai';
import { useState } from 'react';
export default function Chat({
id,
initialMessages,
}: { id?: string | undefined; initialMessages?: UIMessage[] } = {}) {
const [input, setInput] = useState('');
const { sendMessage, messages } = useChat({
id, // use the provided chat ID
messages: initialMessages, // load initial messages
transport: new DefaultChatTransport({
api: '/api/chat',
}),
});
const handleSubmit = (e: React.FormEvent) => {
e.preventDefault();
if (input.trim()) {
sendMessage({ text: input });
setInput('');
}
};
// simplified rendering code, extend as needed:
return (
<div>
{messages.map(m => (
<div key={m.id}>
{m.role === 'user' ? 'User: ' : 'AI: '}
{m.parts
.map(part => (part.type === 'text' ? part.text : ''))
.join('')}
</div>
))}
<form onSubmit={handleSubmit}>
<input
value={input}
onChange={e => setInput(e.target.value)}
placeholder="Type a message..."
/>
<button type="submit">Send</button>
</form>
</div>
);
}存储动作在 toUIMessageStream 的 onEnd 回调中完成——它会收到包含新 AI 响应在内的完整消息列表(UIMessage[]):
import { saveChat } from '@util/chat-store';
import {
convertToModelMessages,
createUIMessageStreamResponse,
streamText,
toUIMessageStream,
UIMessage,
} from 'ai';
import { createGateway } from 'ai';
const gateway = createGateway({ apiKey: process.env.AI_GATEWAY_API_KEY ?? '' });
export async function POST(req: Request) {
const { messages, chatId }: { messages: UIMessage[]; chatId: string } =
await req.json();
const result = streamText({
model: gateway('openai/gpt-5-mini'),
messages: await convertToModelMessages(messages),
});
return createUIMessageStreamResponse({
stream: toUIMessageStream({
stream: result.stream,
originalMessages: messages,
onEnd: ({ messages }) => {
saveChat({ chatId, messages });
},
}),
});
}📌 注意:
useChat的消息格式(面向前端展示,含id、createdAt等字段)不同于ModelMessage格式。官方推荐以useChat格式存储消息。
服务端消息 ID 生成
默认情况下,用户消息 ID 由客户端的 useChat 生成,AI 响应消息 ID 由服务端 streamText 生成。对于持久化场景,应使用在消息入库前就稳定不变的 ID,以保证跨会话一致性、避免恢复消息时的 ID 冲突。
方式一:在 toUIMessageStream 使用 generateMessageId,配合 createIdGenerator() 控制 ID 格式:
import {
createIdGenerator,
createUIMessageStreamResponse,
streamText,
toUIMessageStream,
} from 'ai';
export async function POST(req: Request) {
// ...
const result = streamText({
// ...
});
return createUIMessageStreamResponse({
stream: toUIMessageStream({
stream: result.stream,
originalMessages: messages,
// Generate consistent server-side IDs for persistence:
generateMessageId: createIdGenerator({
prefix: 'msg',
size: 16,
}),
onEnd: ({ messages }) => {
saveChat({ chatId, messages });
},
}),
});
}方式二:用 createUIMessageStream 手动写入带自定义 ID 的 start message part:
import {
generateId,
streamText,
createUIMessageStream,
createUIMessageStreamResponse,
toUIMessageStream,
} from 'ai';
import { createGateway } from 'ai';
const gateway = createGateway({ apiKey: process.env.AI_GATEWAY_API_KEY ?? '' });
export async function POST(req: Request) {
const { messages, chatId } = await req.json();
const stream = createUIMessageStream({
execute: ({ writer }) => {
// Write start message part with custom ID
writer.write({
type: 'start',
messageId: generateId(), // Generate server-side ID for persistence
});
const result = streamText({
model: gateway('openai/gpt-5-mini'),
messages: await convertToModelMessages(messages),
});
writer.merge(
toUIMessageStream({ stream: result.stream, sendStart: false }),
); // omit start message part
},
originalMessages: messages,
onEnd: ({ responseMessage }) => {
// save your chat here
},
});
return createUIMessageStreamResponse({ stream });
}纯客户端应用如果不要求持久化,也可以自定义客户端 ID 生成:
import { createIdGenerator } from 'ai';
import { useChat } from '@ai-sdk/react';
const { ... } = useChat({
generateId: createIdGenerator({
prefix: 'msgc',
size: 16,
}),
// ...
});本章小结
- 用
createUIMessageStream的writer.write可发送三类数据:持久化 data parts(进消息历史)、sources(RAG 引用)、瞬态 parts(仅onData可见); - 同
id写入触发客户端 reconciliation 自动更新,可实现进度条、协作编辑等动态体验; - message metadata 是消息级信息(时间戳/模型/token 用量),经
toUIMessageStream的messageMetadata回调发送,客户端从message.metadata读取; - 持久化三件套:文件/DB 存取
createChat/loadChat/saveChat+ 服务端validateUIMessages校验历史消息 +onEnd回调落盘; - 持久化场景务必使用服务端生成的稳定消息 ID(
generateMessageId或手写 start part)。
🛠️ 动手实践
- 扩展第 16 章的聊天应用:服务端用
createUIMessageStream发送data-status进度部件(loading → success 两态),前端利用 reconciliation 渲染进度指示。 - 给每条 AI 响应附上 message metadata(开始时间 + totalTokens),并在气泡角落显示「耗时模型 X · N tokens」。
- 把聊天记录持久化到本地 JSON 文件:实现
createChat/loadChat/saveChat,加上validateUIMessages校验与服务端generateMessageId,刷新页面后会话可完整恢复。