Skip to content

第 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 的自定义消息类型:

tsx
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 返回:

tsx
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 中:

tsx
writer.write({
  type: 'data-weather',
  id: 'weather-1', // Optional: enables reconciliation
  data: { city: 'San Francisco', status: 'loading' },
});

sources——适合 RAG 场景展示引用的文档或链接(见上方代码第 2 步)。

瞬态(transient)data parts——发送到客户端但不进入消息历史,只能通过 useChatonData 回调访问:

tsx
// 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:

tsx
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:

tsx
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:

tsx
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>;

服务端发送元数据

toUIMessageStreammessageMetadata 回调在不同流式阶段发送元数据:

ts
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 属性读取:

tsx
'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 聊天应用至关重要。本节演示如何结合 useChatstreamText 实现消息持久化。

创建新聊天

用户不带聊天 ID 访问聊天页时,需要先创建新聊天再重定向:

tsx
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 用文件系统存储聊天消息;真实应用应改用数据库或云存储,函数接口设计上可以无缝替换:

tsx
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 目录。

加载已有聊天

tsx
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 校验:

tsx
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 不匹配,要优雅处理校验错误:

tsx
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...
}

展示聊天页面与保存消息

页面组件加载历史后传给聊天组件,后者用 idinitialMessages 初始化 useChat

tsx
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} />;
}
tsx
'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>
  );
}

存储动作在 toUIMessageStreamonEnd 回调中完成——它会收到包含新 AI 响应在内的完整消息列表(UIMessage[]):

tsx
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 的消息格式(面向前端展示,含 idcreatedAt 等字段)不同于 ModelMessage 格式。官方推荐以 useChat 格式存储消息。

服务端消息 ID 生成

默认情况下,用户消息 ID 由客户端的 useChat 生成,AI 响应消息 ID 由服务端 streamText 生成。对于持久化场景,应使用在消息入库前就稳定不变的 ID,以保证跨会话一致性、避免恢复消息时的 ID 冲突。

方式一:在 toUIMessageStream 使用 generateMessageId,配合 createIdGenerator() 控制 ID 格式:

tsx
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:

tsx
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 生成:

tsx
import { createIdGenerator } from 'ai';
import { useChat } from '@ai-sdk/react';

const { ... } = useChat({
  generateId: createIdGenerator({
    prefix: 'msgc',
    size: 16,
  }),
  // ...
});

本章小结

  • createUIMessageStreamwriter.write 可发送三类数据:持久化 data parts(进消息历史)、sources(RAG 引用)、瞬态 parts(仅 onData 可见);
  • id 写入触发客户端 reconciliation 自动更新,可实现进度条、协作编辑等动态体验;
  • message metadata 是消息级信息(时间戳/模型/token 用量),经 toUIMessageStreammessageMetadata 回调发送,客户端从 message.metadata 读取;
  • 持久化三件套:文件/DB 存取 createChat/loadChat/saveChat + 服务端 validateUIMessages 校验历史消息 + onEnd 回调落盘;
  • 持久化场景务必使用服务端生成的稳定消息 ID(generateMessageId 或手写 start part)。

🛠️ 动手实践

  1. 扩展第 16 章的聊天应用:服务端用 createUIMessageStream 发送 data-status 进度部件(loading → success 两态),前端利用 reconciliation 渲染进度指示。
  2. 给每条 AI 响应附上 message metadata(开始时间 + totalTokens),并在气泡角落显示「耗时模型 X · N tokens」。
  3. 把聊天记录持久化到本地 JSON 文件:实现 createChat/loadChat/saveChat,加上 validateUIMessages 校验与服务端 generateMessageId,刷新页面后会话可完整恢复。