跳到主要内容

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 消息返回给用户之前修改或验证消息。必须实现 processOutputResultprocessOutputStream 函数中的一个或两个。

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 完成,再发送响应链中的下一轮请求。