Async Streaming 一等公民设计
Related topics: [[streaming-comparison]]
Overview
Async Streaming 作为一等公民的设计哲学强调:流式不是响应的附加功能,而是交互的核心抽象。这种设计模式下:
- 所有响应都通过流式接口提供,即使是一次性返回的完整内容
- 框架内部以流式为基础构建,非流式只是流式的聚合结果
- 支持细粒度的片段控制、背压和取消机制
- 统一的异步迭代器接口简化消费端代码
Key Concepts
1. 框架对比分析
1.1 kosong (kimi-cli)
流式 API 设计:⭐⭐⭐⭐⭐ 核心抽象
# StreamedMessage 是核心 Protocol
@runtime_checkable
class StreamedMessage(Protocol):
def __aiter__(self) -> AsyncIterator[StreamedMessagePart]:
...
@property
def id(self) -> str | None: ...
@property
def usage(self) -> "TokenUsage | None": ...
- 设计哲学:
generate()始终返回StreamedMessage,流式是唯一接口 - 数据流模式: AsyncIterator (Pull模式) + Callback (Push模式) 混合
- 片段管理:
merge_in_place()方法实现片段合并,支持ContentPart | ToolCall | ToolCallPart - 取消机制: 完整的
asyncio.CancelledError支持,自动清理未完成的 ToolResultFuture - 错误处理: 定义了清晰的错误层次结构 (
ChatProviderError→APIConnectionError/APITimeoutError/APIStatusError)
关键代码片段:
# packages/kosong/src/kosong/_generate.py
async for part in stream:
if on_message_part:
await callback(on_message_part, part.model_copy(deep=True))
if pending_part is None:
pending_part = part
elif not pending_part.merge_in_place(part):
_message_append(message, pending_part)
if isinstance(pending_part, ToolCall) and on_tool_call:
await callback(on_tool_call, pending_part)
pending_part = part
1.2 langchain
流式 API 设计:⭐⭐⭐ 附加功能(但基础架构支持良好)
- 设计哲学:
BaseChatModel同时支持_generate和_stream,流式是可选优化 - 数据流模式: Callback-based (主要) + AsyncIterator (表层)
- 片段管理:
ChatGenerationChunk+merge_chat_generation_chunks()合并 - 取消机制: 通过
disable_streaming标志控制,无原生取消支持 - 错误处理: 错误通过回调传播,支持
run_manager.on_llm_error()
关键代码片段:
# libs/core/langchain_core/language_models/chat_models.py
async for chunk in self._astream(input_messages, stop=stop, **kwargs):
if chunk.message.id is None:
chunk.message.id = run_id
await run_manager.on_llm_new_token(
cast("str", chunk.message.content), chunk=chunk
)
chunks.append(chunk)
yield cast("AIMessageChunk", chunk.message)