Solon v4.1.0

chat - 模型流式事件(详解)

</> markdown
2026年9月7日 上午10:24:45
  • ChatResponse 表示终态,最后的聚合。//在 4.1 之前,也同时表示部分过程态。v4.1 后只表示终态。
  • ChatEvent 表示过程态,不可变的语义事件。//v4.1 后支持

1、从哪里开始

ChatModel 有两条结果路径:

API返回值适用场景
call()一个完整的 ChatResponse不需要观察流式过程
stream()Flux<ChatEvent>打字机、思考展示、工具调用、审计、转发等
// 非流式:直接取得完整结果
ChatResponse response = chatModel.prompt("hello").call();

// 流式:只投影正文增量
chatModel.prompt("hello").stream()
        .filter(e -> e.is(ChatEventType.TEXT_DELTA) && e.hasText())
        .map(ChatEvent::getText)
        .subscribe(System.out::print);

// 流式:取得成功时的完整终态响应
ChatResponse response = chatModel.prompt("hello").stream()
        .filter(e -> e.is(ChatEventType.RESPONSE_END))
        .map(ChatEvent::getResponse)
        .blockFirst();

RESPONSE_ENDgetResponse() 才是完整聚合结果;TEXT_DELTATHINKING_DELTA 等中间事件只表示当前分片。正常完成时一定会产生一个 RESPONSE_END,失败时产生 ERROR 并以同一异常触发 onError,不会再补发 RESPONSE_END。如果订阅方主动取消,Reactor 流会被切断,不能期待取消后再收到事件。

2、常见订阅方式

A. 只显示正文

正文事件必须按 TEXT_DELTA 过滤。不要只按 isDelta() 判断,因为思考、工具参数和拒答也属于增量事件。

public Flux<String> typewriter(String message) {
    return chatModel.prompt(message).stream()
            .filter(e -> e.is(ChatEventType.TEXT_DELTA) && e.hasText())
            .map(ChatEvent::getText);
}

如果同时显示思考和正文,应按分组区分,并只对对应的增量事件读取文本:

chatModel.prompt(query).stream()
        .filter(ChatEvent::isDelta)
        .subscribe(e -> {
            if (e.is(ChatEventType.TEXT_DELTA) && e.hasText()) {
                ui.appendText(e.getText());
            } else if (e.is(ChatEventType.THINKING_DELTA) && e.hasText()) {
                ui.appendThinking(e.getText());
            }
        });

getText() 可能为 null。使用方法引用 map(ChatEvent::getText) 前,应先用 hasText() 排除空负载;Reactor 的 map 不接受 null 结果。

B. 统一消费全部事件

事件类型固定归属于 9 个分组。订阅方可以按 getGroup() 分派,并保留 default,这样新增具体类型时不会因为 switch 漏掉事件。

chatModel.prompt(query).stream().subscribe(
        e -> {
            switch (e.getGroup()) {
                case TEXT:
                    if (e.is(ChatEventType.TEXT_DELTA) && e.hasText()) {
                        ui.appendText(e.getText());
                    }
                    break;
                case THINKING:
                    if (e.is(ChatEventType.THINKING_DELTA) && e.hasText()) {
                        ui.appendThinking(e.getText());
                    }
                    // THINKING_SIGNATURE / THINKING_REDACTED 不是文本增量,按类型单独处理
                    break;
                case TOOL_CALL:
                    handleToolEvent(e);
                    break;
                case SERVER_TOOL:
                    handleServerTool(e.getType(), e.getSubType(), e.getRaw());
                    break;
                case MEDIA:
                    handleMediaOrCitation(e);
                    break;
                case SAFETY:
                    handleSafety(e);
                    break;
                case STEP:
                    handleStep(e.getType(), e.getStep());
                    break;
                case LIFECYCLE:
                    if (e.is(ChatEventType.RESPONSE_END)) {
                        ui.finish(e.getResponse());
                    } else if (e.is(ChatEventType.ABORT)) {
                        ui.abort();
                    }
                    break;
                case META:
                    if (e.is(ChatEventType.USAGE)) {
                        ui.updateUsage(e.getUsage());
                    } else if (e.is(ChatEventType.ERROR)) {
                        ui.fail(e.getError());
                    }
                    break;
                default:
                    log.debug("{} {}", e.getRawType(), e.getRaw());
                    break;
            }
        },
        error -> log.warn("聊天流失败", error));

注意:TEXT_START / TEXT_ENDTHINKING_START / THINKING_END 没有必要读取 getText();它们是边界信号。工具调用的参数也不能当正文处理。

C. 取得终态或异步返回

同步代码可以直接阻塞等待第一个 RESPONSE_END

ChatResponse response = chatModel.prompt(query).stream()
        .filter(e -> e.is(ChatEventType.RESPONSE_END))
        .map(ChatEvent::getResponse)
        .blockFirst();

if (response != null) {
    String text = response.getText();
    List<ToolCall> toolCalls = response.getToolCalls();
}

异步代码保留为 Mono<ChatResponse>,使用 next()

Mono<ChatResponse> response = chatModel.prompt(query).stream()
        .filter(e -> e.is(ChatEventType.RESPONSE_END))
        .map(ChatEvent::getResponse)
        .next();

blockFirst() / next() 都表示“取第一个匹配的终态事件”,不是把中间分片重新聚合。不要对完整事件流直接调用 blockFirst(),否则拿到的通常是 RESPONSE_START

3、事件载荷与类型

分组

分组事件类型主要载荷
LIFECYCLERESPONSE_STARTSTATUSHEARTBEATRESPONSE_ENDABORTRESPONSE_END.getResponse() 为全流终态聚合
STEPSTEP_STARTSTEP_ENDgetStep()STEP_END.getResponse() 为本步聚合,getUsage() 为本步用量
TEXTTEXT_STARTTEXT_DELTATEXT_ENDgetText()
THINKINGTHINKING_STARTTHINKING_DELTATHINKING_ENDTHINKING_SIGNATURETHINKING_REDACTED思考增量或签名使用 getText()
TOOL_CALLTOOL_CALL_STARTTOOL_CALL_ARGS_DELTATOOL_CALL_ENDTOOL_RESULTgetToolCall()getToolCallId()、参数增量 getText()
SERVER_TOOLSERVER_TOOL_STARTSERVER_TOOL_ARGS_DELTASERVER_TOOL_RESULTgetSubType()getRaw()
MEDIACITATIONMEDIA_PARTIALMEDIA_DONEgetBlock()
SAFETYREFUSAL_DELTACONTENT_FILTER拒答文本使用 getText()
METAUSAGEERRORRAWCUSTOMgetUsage()getError()getRaw()

ChatEventType 还提供几个直接的判断方法:

boolean delta = event.isDelta();
boolean terminal = event.isTerminal();
boolean text = event.is(ChatEventType.TEXT_DELTA);
boolean thinking = event.isGroup(ChatEventGroup.THINKING);

isTerminal() 只对 RESPONSE_ENDABORTERROR 返回 trueSTEP_ENDTEXT_ENDTHINKING_ENDTOOL_CALL_END 是各自范围的结束事件,不是全流终止事件。

工具调用的生命周期

客户端工具调用按以下事件序列观察:

void handleToolEvent(ChatEvent e) {
    if (e.is(ChatEventType.TOOL_CALL_START)) {
        onToolStart(e.getToolCallId(), e.getToolCall());
    } else if (e.is(ChatEventType.TOOL_CALL_ARGS_DELTA)) {
        onToolArguments(e.getToolCallId(), e.getText());
    } else if (e.is(ChatEventType.TOOL_CALL_END)) {
        onToolEnd(e.getToolCallId(), e.getToolCall());
    } else if (e.is(ChatEventType.TOOL_RESULT)) {
        onToolResult(e.getToolCallId(), e.getToolCall(), e.getText());
    }
}

TOOL_CALL_ARGS_DELTA 是参数字符串分片,不能假定每一片都是完整 JSON。工具调用完成后,完整的 ToolCall 列表从 STEP_ENDRESPONSE_END 携带的 ChatResponse 读取。服务端工具属于 SERVER_TOOL 分组,不要与本地执行的 TOOL_CALL 混用。

用量与响应

  • USAGE 是某一步的用量事件,getUsage() 可能为空。
  • STEP_ENDgetUsage() 是该步用量。
  • RESPONSE_ENDgetUsage() 是整个 stream()(包括自动工具调用多步)的累计用量。
  • RESPONSE_END / STEP_ENDgetResponse() 是完整聚合;增量事件中的响应(如果有)只是当前帧快照。

4、事件流不变量

核心的 ChatEventNormalizer 会为方言输出补齐边界,但订阅方仍应按以下契约编写:

  1. 正常完成时,全流有一个 RESPONSE_START 和一个 RESPONSE_END;失败时使用 ERROR + onError,不再发 RESPONSE_END
  2. 每轮模型调用有一对 STEP_START / STEP_END;自动工具调用会产生多个步骤。
  3. 每个正文或思考增量都处于对应的 START / END 之间;正文与思考交替时,前一个块会先结束。
  4. 工具参数增量之前会有 TOOL_CALL_START,流结束前会有 TOOL_CALL_END。工具调用采用宽松补齐策略,不应依赖事件去重。
  5. responseId 标识一次 stream()step 标识当前模型调用;providerResponseId 是供应商原始响应 id,可能随自动工具调用的步骤变化。
  6. ABORT 表示服务端或上游中止;调用方主动取消属于 Reactor cancel,两者不是同一个事件。

5、过滤器与失败处理

默认过滤器只屏蔽 HEARTBEATRAWLIFECYCLESTEP 分组始终放行,因为它们承载流的生命周期和终态聚合。

// 默认事件之外,再放行未建模的 RAW 事件
Flux<ChatEvent> events = chatModel.prompt(query)
        .eventFilter(ChatEventFilter.DEFAULT.or(
                ChatEventFilter.of(ChatEventType.RAW)))
        .stream();

// 需要全部事件时显式开启(包括 HEARTBEAT 与 RAW)
Flux<ChatEvent> allEvents = chatModel.prompt(query)
        .eventFilter(ChatEventFilter.all())
        .stream();

失败有两个观察点,但表示同一次失败:

chatModel.prompt(query).stream().subscribe(
        event -> {
            if (event.is(ChatEventType.ERROR)) {
                // 可选:保存已经完成的部分结果
                savePartial(event.getResponse(), event.getUsage());
            }
        },
        error -> retryOrRecover(error));

需要展示或保存已完成部分时消费 ERROR;需要使用 retryWhenonErrorResume 等 Reactor 操作符时处理 onError。不要同时把两条通道当作两次独立失败。

6、方言作者:解析入口

ChatDialect 的解析入口是:

void parseResponseJson(ChatStreamContext ctx, String respJson);

方言可以:

  • 把正文、思考和工具调用写入 ctx.getAccumulator() 的内容项,由核心统一转换成 TEXT_*THINKING_*TOOL_CALL_*
  • 对生命周期、服务端工具、引用、拒答、思考签名等扩展语义,使用 ctx.emit(...) 发射事件;
  • 使用 ctx.event(type) 创建事件。该构建器已经预填当前 responseId、供应商响应 id(如果已设置)和 step
  • 使用 ctx.attrPut / ctx.attrAs 保存跨帧的方言私有状态;
  • 解析错误时写入 ctx.getAccumulator().setError(...),已消费但没有语义内容时不发事件。

示例:

@Override
public void parseResponseJson(ChatStreamContext ctx, String json) {
    // 供应商首帧提取到原始响应 id 后记录一次
    ctx.setProviderResponseId(readProviderResponseId(json));

    String text = readTextDelta(json);
    if (text != null && !text.isEmpty()) {
        ctx.getAccumulator().addContentItem(new AssistantMessage(text));
    }

    String citationUrl = readCitation(json);
    if (citationUrl != null) {
        ctx.emit(ctx.event(ChatEventType.CITATION)
                .block(buildCitation(citationUrl))
                .build());
    }
}

方言不要同时把同一份正文、思考或工具调用既写入内容项又通过 ctx.emit(...) 发射。内容主干应选择内容项这一条路径,否则订阅方可能收到重复增量,终态聚合也会重复。ChatEventNormalizer 可以为第三方方言补齐缺失边界,但兜底不应替代正常的解析设计。

7、从旧流式用法迁移

旧思路4.1 写法
Flux<ChatResponse> 的每个对象当作累计结果消费 Flux<ChatEvent>;完整结果从 RESPONSE_END.getResponse() 读取
用响应对象的可变字段拼接正文过滤 TEXT_DELTA,读取 getText();最终结果读取 ChatResponse.getText()
blockLast() 等待一个可变末帧过滤 RESPONSE_END 后使用 blockFirst(),或异步使用 next()
用启发式判断区分正文、思考和工具根据 ChatEventTypeChatEventGroup 分派
自己维护正文/思考的开始和结束状态依赖 *_START / *_DELTA / *_END 事件边界
从每个参数分片解析完整工具调用累积 TOOL_CALL_ARGS_DELTA;完整工具调用从终态 ChatResponse.getToolCalls() 读取
方言保留旧的布尔返回解析入口实现 parseResponseJson(ChatStreamContext, String)