流式查询

Explore/Data Asking 不是普通的请求—响应接口。analyze_data_stream(Python)和 AnalyzeDataStream(Go)返回 Server-Sent Events(SSE)流;应用必须持续读取事件、保存初始化事件中的 request_id,并在结束或异常时关闭流。

请求

两套 SDK 的请求包含相同字段:

JSON 字段

是否必需

作用

question

本次分析的问题;SDK 拒绝空字符串

source

调用来源标识

session_id

继续已有会话

session_name

会话名称

config

数据源、范围、过滤和上下文配置

数据表与文件的选择方式参见 数据源

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:源码注释列出 ragnl2sql

  • data:随事件变化的对象;

  • step_typestep_name:部分 NL2SQL 或步骤事件使用;

  • raw_data / RawData:原始 JSON,供兼容未知事件或排障。

SDK 源码列出的事件包括:

事件

用途

init

首个初始化事件;data 中包含 request_idsession_title

classification

问题分类

decomposition

归因类问题的拆解

step_start / step_complete

归因步骤开始与完成

chunks / answer_chunk

RAG 数据和增量回答

step_type / step_name

NL2SQL 过程数据

complete

分析完成

error

流内错误

不要假设所有事件都通过 type 标识,也不要把 data 反序列化成一个固定结构。先检查 typesourcestep_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_idstatususer_iduser_name。本地取消上下文或关闭网络连接不等同于服务端取消;需要停止后端运行时调用取消方法。

错误分层

阶段

典型失败

处理

建立流之前

认证、权限、空问题、HTTP 错误

不进入读取循环;修正请求后再试

读取过程中

网络中断、读取超时、SSE 格式问题

关闭流并记录最近事件和 request_id

流内 error 事件

检索、SQL 或分析步骤失败

展示服务错误;不要把已收到的片段标记为完整回答

complete 前 EOF

连接正常结束但未见完成事件

视为不完整运行并保留诊断信息

流式分析可能已经执行了部分工作。自动重试前应判断操作是否安全,并创建新的运行记录;不要把两次流的片段拼成一个回答。

会话与持久化

传入 session_id 可继续会话,session_name 可提供名称。应用自己的记录至少保存:

  • 会话 ID;

  • 初始化事件返回的请求 ID;

  • 问题和数据源配置摘要;

  • 按顺序到达的事件或必要审计字段;

  • 是否收到 complete、是否取消及最终错误。

保存来源和完成状态比只保存渲染后的文本更利于复现问题。