Agent.stream()
.stream() 方法支持从 Agent 实时流式传输响应,并提供增强功能和灵活的格式。该方法接受消息和可选的流式传输选项,可提供现代化的流式传输体验,同时兼容 Mastra 原生格式和 AI SDK v5+。
使用示例使用示例的直接链接
const stream = await agent.stream('message for agent')
信息
模型兼容性:此方法专为 V2 模型设计。V1 模型应使用 .streamLegacy() 方法。框架会自动检测模型版本,并在版本不匹配时抛出错误。
参数参数的直接链接
messages:
string | string[] | CoreMessage[] | AiMessageType[] | UIMessageWithMetadata[]
要发送给 Agent 的消息。可以是单个字符串、字符串数组或结构化消息对象。
options?:
AgentExecutionOptions<Output, Format>
流式传输过程的可选配置。
AgentExecutionOptions<Output, Format>
maxSteps?:
number
执行期间最多运行的步骤数。
scorers?:
MastraScorers | Record<string, { scorer: MastraScorer['name']; sampling?: ScoringSamplingConfig }>
对执行结果运行的评估 scorer。
scorer:
string
要使用的 scorer 名称。
sampling?:
ScoringSamplingConfig
scorer 的采样配置。
type:
'none' | 'ratio'
采样策略的类型。使用 'none' 禁用采样,或使用 'ratio' 进行按比例采样。
rate?:
number
采样率(0-1)。当类型为 'ratio' 时必填。
onIterationComplete?:
(context: IterationCompleteContext) => { continue?: boolean; feedback?: string } | void | Promise<{ continue?: boolean; feedback?: string } | void>
每次迭代完成后调用的回调函数。可用于监控进度、提供反馈来引导 Agent,或提前停止执行。该回调接收迭代上下文,包括当前文本、Tool 调用和结束原因。
context.iteration:
number
当前迭代次数(从 1 开始)。
context.maxIterations:
number | undefined
允许的最大迭代次数(如果已设置)。
context.text:
string
本次迭代的文本响应。
context.isFinal:
boolean
这是否为最后一次迭代。
context.finishReason:
string
本次迭代结束的原因(例如 'stop'、'length'、'tool-calls')。
context.toolCalls:
ToolCall[]
本次迭代中发起的 Tool 调用。
context.messages:
MastraDBMessage[]
目前累计的所有消息。
return.continue?:
boolean
设为 false 可提前停止执行。
return.feedback?:
string
用于引导 Agent 下一次迭代的反馈消息。
isTaskComplete?:
IsTaskCompleteConfig
用于验证任务是否完成的任务完成度评分配置。使用 Mastra 的评估 scorer 自动检查 Agent 的响应是否满足完成标准。
scorers:
MastraScorer[]
用于评估任务完成情况的 scorer 数组。每个 scorer 返回 0(失败)或 1(通过)。
strategy?:
'all' | 'any'
组合 scorer 结果的策略。'all' 要求所有 scorer 均通过,'any' 要求至少一个通过。
onComplete?:
(result: IsTaskCompleteRunResult) => void | Promise<void>
任务完成检查结束时调用的回调。接收包含各个 scorer 分数的结果。
parallel?:
boolean
是否并行运行 scorer。
timeout?:
number
等待所有 scorer 完成的最长时间(毫秒)。
suppressFeedback?:
boolean
为 true 时,标记完成检查反馈,使使用方可以在显示的输出中将其隐藏。仅当检查失败时,反馈才会添加到对话中,以引导下一次迭代。
delegation?:
DelegationConfig
Subagent 委派的配置。可用于控制和监控 Agent 何时将任务委派给其他 Agent,包括修改、拒绝委派以及提供反馈来引导 supervisor。
onDelegationStart?:
(context: DelegationStartContext) => DelegationStartResult | void | Promise<DelegationStartResult | void>
委派给 Subagent 前调用。可用于修改委派参数、完全拒绝委派,或更改
context.requestContext,向 Subagent 运行的 request context 添加条目。onDelegationComplete?:
(context: DelegationCompleteContext) => { feedback?: string } | void | Promise<{ feedback?: string } | void>
Subagent 委派完成后调用。上下文包含用于停止后续执行的
bail() 方法,你可以返回 { feedback } 来引导 supervisor 的下一步操作。反馈会作为 assistant 消息保存到 supervisor memory 中。messageFilter?:
(context: MessageFilterContext) => MastraDBMessage[] | Promise<MastraDBMessage[]>
委派给 Subagent 前调用的回调函数。可用于筛选传递给 Subagent 的消息。
tracingContext?:
TracingContext
用于 span 层级和元数据的 Tracing 上下文。
returnScorerData?:
boolean
是否在响应中返回详细的评分数据。
onChunk?:
(chunk: ChunkType) => Promise<void> | void
流式传输期间针对每个数据块调用的回调函数。
onError?:
({ error }: { error: Error | string }) => Promise<void> | void
流式传输期间发生错误时调用的回调函数。
onAbort?:
(event: any) => Promise<void> | void
Stream 中止时调用的回调函数。
abortSignal?:
AbortSignal
用于中止 Agent 执行的信号对象。信号中止时,所有正在进行的操作都会终止,包括 Agent 委派且仍在运行的任何 Subagent。
activeTools?:
Array<keyof ToolSet> | undefined
执行期间可以使用的活跃 Tool 名称数组。
prepareStep?:
PrepareStepFunction<any>
多步骤执行中每个步骤之前调用的回调函数。
context?:
ModelMessage[]
提供给 Agent 的额外上下文消息。
structuredOutput?:
StructuredOutputOptions<S extends ZodTypeAny = ZodTypeAny>
用于微调结构化输出生成的选项。
schema:
StandardJSONSchemaV1
定义预期输出结构的标准 JSON Schema。
model?:
MastraLanguageModel
用于生成结构化输出的语言模型。提供后,Agent 可在多步骤响应中包含 Tool 调用、文本和结构化输出
errorStrategy?:
'strict' | 'warn' | 'fallback'
处理 schema 验证错误的策略。'strict' 抛出错误,'warn' 记录警告,'fallback' 使用回退值。
fallbackValue?:
<S extends ZodTypeAny>
schema 验证失败且 errorStrategy 为 'fallback' 时使用的回退值。
instructions?:
string
为结构化输出模型提供的额外指令。
jsonPromptInjection?:
boolean | 'system' | 'inline' | 'auto'
控制 JSON schema 如何传递给模型。设为 'auto' 时,在支持的情况下使用原生结构化输出,否则使用内联提示词注入。
providerOptions?:
ProviderOptions
传递给内部结构化 Agent 的 Provider 专用选项。可用于控制模型行为,例如思考模型的推理强度(如
{ openai: { reasoningEffort: 'low' } })。outputProcessors?:
Processor[]
覆盖 Agent 上设置的输出 Processor。输出 Processor 可在 Agent 消息返回给用户之前修改或验证消息。必须实现
processOutputResult 和 processOutputStream 函数中的一个或两个。includeRawChunks?:
boolean
是否在 Stream 输出中包含原始数据块(并非所有模型 Provider 都支持)。
inputProcessors?:
Processor[]
覆盖 Agent 上设置的输入 Processor。输入 Processor 可在消息由 Agent 处理前修改或验证消息。必须实现
processInput 函数。instructions?:
string
针对本次生成覆盖 Agent 默认指令的自定义指令。无需创建新的 Agent 实例,即可动态修改 Agent 行为。
system?:
string | string[] | CoreSystemMessage | SystemModelMessage | CoreSystemMessage[] | SystemModelMessage[]
要包含在提示词中的自定义 system 消息。可以是单个字符串、消息对象或两者之一的数组。system 消息提供额外的上下文或行为指令,用于补充 Agent 的主要指令。
output?:
Zod schema | JsonSchema7
**已弃用。** 使用不含 model 的 structuredOutput 可实现相同效果。定义预期输出结构。可以是 JSON Schema 对象或 Zod schema。
memory?:
object
Memory 配置。这是管理 Memory 的首选方式。
thread:
string | { id: string; metadata?: Record<string, any>, title?: string }
对话 thread,可以是字符串 ID,也可以是包含
id 和可选 metadata 的对象。resource:
string
与 thread 关联的用户或资源标识符。
options?:
MemoryConfig
Memory 行为配置,包括 lastMessages、readOnly、semanticRecall、workingMemory 和 filterIncompleteToolCalls。
onTitleGenerated?:
(title: string) => void | Promise<void>
生成 thread 标题并将其持久化到存储后异步触发的回调。标题生成在后台运行,可能会在 Stream 结束后完成。仅当 Memory 选项中启用了
generateTitle 且 thread 没有现有标题时触发。onFinish?:
StreamTextOnFinishCallback<any> | StreamObjectOnFinishCallback<OUTPUT>
流式传输完成时调用的回调函数。接收最终结果。
onStepFinish?:
StreamTextOnStepFinishCallback<any> | never
每个执行步骤完成后调用的回调函数。以 JSON 字符串形式接收步骤详情。结构化输出不支持此选项
telemetry?:
TelemetrySettings
流式传输期间的 OTLP telemetry 收集设置(不是 Tracing)。
isEnabled?:
boolean
启用或禁用 telemetry。实验阶段默认禁用。
recordInputs?:
boolean
启用或禁用输入记录。默认启用。为避免记录敏感信息,你可能需要禁用输入记录。
recordOutputs?:
boolean
启用或禁用输出记录。默认启用。为避免记录敏感信息,你可能需要禁用输出记录。
functionId?:
string
此函数的标识符。用于按函数对 telemetry 数据分组。
modelSettings?:
CallSettings
Model-specific settings like temperature, maxOutputTokens, topP, etc. These settings control how the language model generates responses.
temperature?:
number
Controls randomness in generation (0-2). Higher values make output more random.
maxOutputTokens?:
number
Maximum number of tokens to generate in the response. Note: Use maxOutputTokens (not maxTokens) as per AI SDK v5 convention.
maxRetries?:
number
Maximum number of retry attempts for failed requests.
topP?:
number
Nucleus sampling parameter (0-1). Controls diversity of generated text.
topK?:
number
Top-k sampling parameter. Limits vocabulary to k most likely tokens.
presencePenalty?:
number
Penalty for token presence (-2 to 2). Reduces repetition.
frequencyPenalty?:
number
Penalty for token frequency (-2 to 2). Reduces repetition of frequent tokens.
stopSequences?:
string[]
Stop sequences. If set, the model will stop generating text when one of the stop sequences is generated.
toolChoice?:
'auto' | 'none' | 'required' | { type: 'tool'; toolName: string }
控制 Agent 在流式传输期间如何使用 Tool。
'auto':
string
让模型决定是否使用 Tool(默认)。
'none':
string
不使用任何 Tool。
'required':
string
要求模型至少使用一个 Tool。
{ type: 'tool'; toolName: string }:
object
要求模型按名称使用特定 Tool。
toolsets?:
ToolsetsInput
流式传输期间额外提供给 Agent 的 Toolset。
clientTools?:
ToolsInput
在请求的 'client' 端执行的 Tool。这些 Tool 的定义中没有 execute 函数。
hooks?:
ToolHooks
在 Tool 调用前后运行的单次执行 hook。会覆盖本次执行中匹配的 Agent 级 hook。
beforeToolCall 可以返回 { proceed: false, output } 以跳过 Tool 调用。savePerStep?:
boolean
每个 Stream 步骤完成后增量保存消息(默认:false)。
requireToolApproval?:
boolean
为 true 时,所有 Tool 调用在执行前都需要显式批准。Stream 会发出
tool-call-approval 数据块并暂停,直到调用 approveToolCall() 或 declineToolCall()。autoResumeSuspendedTools?:
boolean
为 true 时,当用户在同一 thread 中发送新消息,会自动恢复已暂停的 Tool。Agent 会根据 Tool 的
resumeSchema 从用户消息中提取 resumeData。需要配置 Memory。toolCallConcurrency?:
number
并发执行的 Tool 调用最大数量。可能需要批准时默认为 1,否则为 10。
providerOptions?:
Record<string, Record<string, JSONValue>>
传递给底层 LLM Provider 的额外 Provider 专用选项。结构为
{ providerName: { optionKey: value } }。例如:{ openai: { reasoningEffort: 'high' }, anthropic: { maxTokens: 1000 } }。openai?:
Record<string, JSONValue>
OpenAI 专用选项。示例:
{ reasoningEffort: 'high' }anthropic?:
Record<string, JSONValue>
Anthropic 专用选项。示例:
{ maxTokens: 1000 }google?:
Record<string, JSONValue>
Google 专用选项。示例:
{ safetySettings: [...] }[providerName]?:
Record<string, JSONValue>
其他 Provider 专用选项。键为 Provider 名称,值为 Provider 专用选项的记录。
runId?:
string
本次生成运行的唯一 ID。可用于跟踪和调试。
requestContext?:
RequestContext
用于依赖注入和上下文信息的 Request Context。
tracingContext?:
TracingContext
用于创建子 span 和添加元数据的 Tracing 上下文。使用 Mastra 的 Tracing 系统时会自动注入。
currentSpan?:
Span
用于创建子 span 和添加元数据的当前 span。可用于在执行期间创建自定义子 span 或更新 span 属性。
tracingOptions?:
TracingOptions
Tracing 配置选项。
metadata?:
Record<string, any>
添加到根 Trace span 的元数据。可用于添加用户 ID、会话 ID 或功能标志等自定义属性。
requestContextKeys?:
string[]
要提取为此 Trace 元数据的额外 RequestContext 键。支持使用点表示法访问嵌套值(例如 'user.id')。
traceId?:
string
本次执行使用的 Trace ID(1-32 个十六进制字符)。如果提供,此 Trace 将成为指定 Trace 的一部分。
parentSpanId?:
string
本次执行使用的父 span ID(1-16 个十六进制字符)。如果提供,根 span 将创建为此 span 的子 span。
tags?:
string[]
应用于此 Trace 的标签。用于分类和筛选 Trace 的字符串标签。
versions?:
VersionOverrides
针对每次调用的 Subagent 委派版本覆盖。该配置会合并到 Mastra 实例级版本之上,并通过 requestContext 自动传播到 Subagent 调用。需要 editor 包。请参阅 Editor 版本控制。
VersionOverrides
agents?:
Record<string, VersionSelector>
Agent ID 到其版本选择器的映射。
VersionSelector
versionId?:
string
按 ID 指定特定版本。
status?:
'draft' | 'published'
指定具有此发布状态的最新版本。
untilIdle?:
boolean | { maxIdleMs?: number }
设置后,Stream 会在后台任务继续执行期间保持打开。后台任务完成时,Agent 会自动重新调用 LLM,并通过同一个
fullStream 流式传输后续轮次。传入 true 使用默认设置(5 分钟空闲超时),或传入包含 maxIdleMs 的对象进行配置。需要 Memory。替代独立的 streamUntilIdle() 方法。maxIdleMs?:
number
轮次之间空闲达到指定毫秒数后关闭外层 Stream。仅当包装器处于轮次之间时,计时器才会运行。默认值:5 分钟。
返回值返回值的直接链接
stream:
MastraModelOutput<Output>
返回一个 MastraModelOutput 实例,用于访问流式输出。
traceId?:
string
启用 Tracing 时,与本次执行关联的 Trace ID。可用于关联日志和调试执行流程。
spanId?:
string
启用 Tracing 时,与本次执行关联的根 span ID。可用于 span 级查询和关联。
扩展使用示例扩展使用示例的直接链接
Mastra 格式(默认)Mastra 格式(默认)的直接链接
index.ts
import { stepCountIs } from 'ai-v5'
const stream = await agent.stream('Tell me a story', {
stopWhen: stepCountIs(3), // Stop after 3 steps
modelSettings: {
temperature: 0.7,
},
})
// Access text stream
for await (const chunk of stream.textStream) {
console.log(chunk)
}
// or access full stream
for await (const chunk of stream.fullStream) {
console.log(chunk)
}
// Get full text after streaming
const fullText = await stream.text
AI SDK v5+ 格式AI SDK v5+ 格式的直接链接
要在 AI SDK v5(及更高版本)中使用此 Stream,可以通过工具函数 toAISdkStream 进行转换。
index.ts
import { stepCountIs, createUIMessageStreamResponse } from 'ai'
import { toAISdkStream } from '@mastra/ai-sdk'
const stream = await agent.stream('Tell me a story', {
stopWhen: stepCountIs(3), // Stop after 3 steps
modelSettings: {
temperature: 0.7,
},
})
// In an API route for frontend integration
return createUIMessageStreamResponse({
stream: toAISdkStream(stream, { from: 'agent' }),
})
使用回调使用回调的直接链接
所有回调函数现在都可以作为顶层属性使用,使 API 使用体验更加简洁。
index.ts
const stream = await agent.stream('Tell me a story', {
onFinish: result => {
console.log('Streaming finished:', result)
},
onStepFinish: step => {
console.log('Step completed:', step)
},
onChunk: chunk => {
console.log('Received chunk:', chunk)
},
onError: ({ error }) => {
console.error('Streaming error:', error)
},
onAbort: event => {
console.log('Stream aborted:', event)
},
})
// Process the stream
for await (const chunk of stream.textStream) {
console.log(chunk)
}
使用选项的高级示例使用选项的高级示例的直接链接
index.ts
import { z } from 'zod'
import { stepCountIs } from 'ai'
await agent.stream('message for agent', {
stopWhen: stepCountIs(3), // Stop after 3 steps
modelSettings: {
temperature: 0.7,
},
memory: {
thread: 'user-123',
resource: 'test-app',
},
toolChoice: 'auto',
// Structured output with better DX
structuredOutput: {
schema: z.object({
sentiment: z.enum(['positive', 'negative', 'neutral']),
confidence: z.number(),
}),
model: 'openai/gpt-5.6-sol',
errorStrategy: 'warn',
},
// Output processors for streaming response validation
outputProcessors: [
new ModerationProcessor({ model: 'openrouter/openai/gpt-oss-safeguard-20b' }),
new BatchPartsProcessor({ maxBatchSize: 3, maxWaitTime: 100 }),
],
})
Responses WebSocket 传输Responses WebSocket 传输的直接链接
通过 Provider 选项启用 Responses WebSocket 流式传输。此选项仅适用于流式调用,并支持 OpenAI 直连模型和 Azure OpenAI Responses 部署。如果 WebSocket 流式传输不可用,Mastra 会回退到 HTTP 流式传输。默认情况下,Stream 完成时 Mastra 会关闭 WebSocket。
index.ts
const stream = await agent.stream('Hello', {
providerOptions: {
openai: {
transport: 'websocket', // 'websocket' | 'fetch' | 'auto'
websocket: {
url: 'wss://api.openai.com/v1/responses',
closeOnFinish: true, // default
},
},
},
})
对于 Azure OpenAI,请使用 useResponsesAPI: true 配置 gateway,然后使用 providerOptions.azure.transport。
index.ts
const stream = await agent.stream('Hello', {
providerOptions: {
azure: {
transport: 'websocket',
store: false,
websocket: { closeOnFinish: true },
},
},
})
要在 Stream 完成后保持连接打开,请设置 closeOnFinish: false,并手动将其关闭。
index.ts
const stream = await agent.stream('Hello', {
providerOptions: {
openai: {
transport: 'websocket',
websocket: { closeOnFinish: false },
},
},
})
// Later, when you're done with the connection:
stream.transport?.close()
Responses WebSocket 连接一次运行一个响应。对于同一 WebSocket transport,Mastra 会拒绝包含 previous_response_id 的重叠后续请求。请等待当前 Stream 完成,再发送响应链中的下一轮请求。