流式查询¶
Explore/Data Asking 不是普通的请求—响应接口。analyze_data_stream(Python)和 AnalyzeDataStream(Go)返回 Server-Sent Events(SSE)流;应用必须持续读取事件、保存初始化事件中的 request_id,并在结束或异常时关闭流。
请求¶
两套 SDK 的请求包含相同字段:
JSON 字段 |
是否必需 |
作用 |
|---|---|---|
|
是 |
本次分析的问题;SDK 拒绝空字符串 |
|
否 |
调用来源标识 |
|
否 |
继续已有会话 |
|
否 |
会话名称 |
|
否 |
数据源、范围、过滤和上下文配置 |
数据表与文件的选择方式参见 数据源。
Python:完整读取事件¶
from moi import RawClient
raw = RawClient("https://api.example.com", "your-api-key")
request = {
"question": "Summarize the renewal conditions.",
"session_id": session_id,
"config": {
"data_source": {
"type": "specified",
"files": {
"type": "specified",
"file_id_list": selected_file_ids,
},
}
},
}
request_id = None
with raw.analyze_data_stream(request) as stream:
while True:
event = stream.read_event()
if event is None:
break
if event.step_type == "init":
init = event.get_init_event_data()
if init is not None:
request_id = init.request_id
# event.type, event.source, event.step_type, event.step_name
# 和 event.data 用于按事件类别更新界面或运行记录。
handle_event(event)
read_event() 在流结束时返回 None。上下文管理器会调用 close();如果不使用 with,必须在 finally 中关闭。
Go:完整读取事件¶
package main
import (
"context"
"io"
sdk "github.com/matrixorigin/moi-go-sdk"
)
func analyze(ctx context.Context, client *sdk.RawClient, fileIDs []string) error {
req := &sdk.DataAnalysisRequest{
Question: "Summarize the renewal conditions.",
Config: &sdk.DataAnalysisConfig{
DataSource: &sdk.DataSource{
Type: "specified",
Files: &sdk.FileConfig{
Type: "specified",
FileIDList: fileIDs,
},
},
},
}
stream, err := client.AnalyzeDataStream(ctx, req)
if err != nil {
return err
}
defer stream.Close()
for {
event, err := stream.ReadEvent()
if err == io.EOF {
return nil
}
if err != nil {
return err
}
if event.StepType == "init" {
init := event.GetInitEventData()
if init != nil {
saveRequestID(init.RequestID)
}
}
handleEvent(event)
}
}
Go 的 ReadEvent() 以 io.EOF 表示正常结束。AnalyzeDataStream 接受 CallOption;大事件可使用 WithStreamBufferSize 调整初始缓冲区。读取超时按“相邻消息之间的等待时间”计算,不是整个分析的总时长。
事件模型¶
公开 DataAnalysisStreamEvent 的稳定外层字段是:
type:例如分类、完成或错误类事件;source:源码注释列出rag和nl2sql;data:随事件变化的对象;step_type、step_name:部分 NL2SQL 或步骤事件使用;raw_data/RawData:原始 JSON,供兼容未知事件或排障。
SDK 源码列出的事件包括:
事件 |
用途 |
|---|---|
|
首个初始化事件; |
|
问题分类 |
|
归因类问题的拆解 |
|
归因步骤开始与完成 |
|
RAG 数据和增量回答 |
|
NL2SQL 过程数据 |
|
分析完成 |
|
流内错误 |
不要假设所有事件都通过 type 标识,也不要把 data 反序列化成一个固定结构。先检查 type、source 和 step_type,未知事件记录 raw_data 后忽略或转交兼容处理。
取消运行¶
初始化事件给出取消所需的 request_id。只有发起该请求的用户可以取消。
Python:
result = raw.cancel_analyze({"request_id": request_id})
Go:
result, err := client.CancelAnalyze(ctx, &sdk.CancelAnalyzeRequest{
RequestID: requestID,
})
取消响应包含 request_id、status、user_id 和 user_name。本地取消上下文或关闭网络连接不等同于服务端取消;需要停止后端运行时调用取消方法。
错误分层¶
阶段 |
典型失败 |
处理 |
|---|---|---|
建立流之前 |
认证、权限、空问题、HTTP 错误 |
不进入读取循环;修正请求后再试 |
读取过程中 |
网络中断、读取超时、SSE 格式问题 |
关闭流并记录最近事件和 |
流内 |
检索、SQL 或分析步骤失败 |
展示服务错误;不要把已收到的片段标记为完整回答 |
|
连接正常结束但未见完成事件 |
视为不完整运行并保留诊断信息 |
流式分析可能已经执行了部分工作。自动重试前应判断操作是否安全,并创建新的运行记录;不要把两次流的片段拼成一个回答。
会话与持久化¶
传入 session_id 可继续会话,session_name 可提供名称。应用自己的记录至少保存:
会话 ID;
初始化事件返回的请求 ID;
问题和数据源配置摘要;
按顺序到达的事件或必要审计字段;
是否收到
complete、是否取消及最终错误。
保存来源和完成状态比只保存渲染后的文本更利于复现问题。