Day 17 - 打開 streaming

Day 17 - 打開 streaming

約 14,199 字

第 2 天攔下來的那包資料裡有一個 "stream": true,這個設定是要 API 把回應一段一段先送回來,不用等模型整段寫完才一次給。KeSi 現在還沒開這個設定,所以每次回話都是按下 Enter 之後畫面停住,等模型整段想完文字才一次全部出現。今天就來把串流打開,順便處理一個打開之後才會遇到的麻煩,就是回答印到一半的時候按下 Ctrl-C 中斷的話,日記要怎麼辦?

GitHub Repo:https://github.com/kaochenlong/KeSi

串流送什麼回來?

如果打開了串流,API 送回來的不再是一整包回應而是一連串的事件。為了看清楚這些事件長什麼樣子,我寫了一支小程式 examples/day17/events.py,它會用 KeSi 平常送出去的 system prompt 跟工具清單,問一句會用到工具的「幫我讀一下 note.txt 的內容」,然後把收到的每一個事件跟收到的時間印出來。這支程式只看事件、不會真的執行工具:

uv run --env-file .env examples/day17/events.py
# 1. 原始事件:幫我讀一下 note.txt 的內容
​
  [ 1.05s] message_start
  [ 1.23s] content_block_start type=text
  [ 1.23s] text_delta       '我'
  [ 1.23s] text_delta       '幫你'
  [ 1.23s] text_delta       '讀一'
  [ 1.23s] text_delta       '下 note.txt '
  [ 1.23s] text_delta       '的內容。'
  [ 1.28s] content_block_stop
  [ 1.28s] content_block_start type=tool_use
  [ 1.28s] input_json_delta ''
  [ 1.46s] input_json_delta '{"file_'
  [ 1.46s] input_json_delta 'path": '
  [ 1.46s] input_json_delta '"note.t'
  [ 1.46s] input_json_delta 'xt"}'
  [ 1.46s] content_block_stop
  [ 1.46s] message_delta stop_reason=tool_use output_tokens=76
  [ 1.46s] message_stop
​
  迭代時拿到的事件種類與次數:
    content_block_delta  10
    content_block_start  2
    content_block_stop   2
    input_json           5
    message_delta        1
    message_start        1
    message_stop         1
    text                 5

第 3 天說模型的回覆是「一疊積木」,這裡看到的就是那疊積木一塊一塊送過來的樣子。每一塊積木從 content_block_start 開始,中間是一段一段的 delta,delta 就是這次新增的那一小段內容,最後用 content_block_stop 收尾,文字積木的 delta 是 text_delta,工具積木的是 input_json_delta。整則訊息的頭尾是 message_start 跟 message_stop,中間的 message_delta 帶著 stop_reason 跟 output_tokens。

工具那一塊的參數不是一次給的,而是一段一段拼起來的 JSON。第一段是空字串,跟著 content_block_start 一起到,後面四段是 {"file_、path":、"note.t、xt"},切的位置完全不管 JSON 的語法,"note.t 連字串都還沒結束。

收到的事件也不只上面這幾種。官方的串流文件提到串流裡可能穿插任意數量的 ping 事件,途中也可能收到 error 事件,例如服務忙碌時的 overloaded_error,另外還有一句:

In accordance with the versioning policy, new event types may be added, and your code should handle unknown event types gracefully.

這句話的意思是,以後可能會多出新的事件種類,程式碰到不認得的事件不能因此出錯。KeSi 的程式只會看三種事件,例如收到 text_delta 就把那一小段文字印到畫面上,收到 content_block_stop 就記下又有一塊積木完整了,收到 message_start 就記下回應已經開始,這兩個記錄是給中途按 Ctrl-C 的時候用的,待會後面會講到。其他像 ping 或是以後才新增的事件,收到了也不理它,所以就算哪天多了新的事件種類,程式也不會因此而壞掉。

輸出最後那張次數表還有一個小地方,text 跟 input_json 這兩種事件不是 API 送的,是 SDK 看到 text_delta 跟 input_json_delta 之後自己另外產生的事件,方便只想拿文字或工具參數的人直接使用,所以各有 5 次,加起來剛好是 content_block_delta 的 10 次。

邊走邊吃、邊收邊印

串流本身的部分其實不長,client.messages.create() 換成 client.messages.stream() 就好。我把送出請求的這一段抽成 stream_turn() 函式,最上面的 StreamFailure 跟後面兩個 except 是替按下 Ctrl-C 跟串流途中出錯準備的:

class StreamFailure(Exception):
    """串流已開始後發生的 API 錯誤,連同當下可用的快照一起帶出去。"""
​
    def __init__(self, cause, snapshot, finished):
        super().__init__(str(cause))
        self.cause = cause
        self.snapshot = snapshot
        self.finished = finished
​
​
def stream_turn(client, history):
    """送出一輪並分段印出。回傳(回應、完成的積木數、是否被中斷)。
​
    被中斷時回的是快照,而快照的最後一塊可能組到一半,所以要一起回報
    「有幾塊是完整的」,呼叫端才知道哪些能放進日記。
    """
    stream = None
    started = False
    finished = 0
    printed = False
    try:
        with client.messages.stream(
            model="claude-haiku-4-5",
            max_tokens=1024,
            system=system_prompt(),
            tools=TOOLS,
            messages=history,
        ) as active:
            stream = active
            for event in active:
                if event.type == "message_start":
                    started = True
                elif event.type == "content_block_stop":
                    finished = event.index + 1
                elif (
                    event.type == "content_block_delta"
                    and event.delta.type == "text_delta"
                ):
                    if not printed:
                        print("KeSi > ", end="", flush=True)
                        printed = True
                    print(event.delta.text, end="", flush=True)
            if printed:
                print()
            return active.get_final_message(), finished, False
    except KeyboardInterrupt:
        if printed:
            print()
        snapshot = (
            stream.current_message_snapshot
            if stream is not None and started
            else None
        )
        return snapshot, finished, True
    except anthropic.AnthropicError as exc:
        if stream is None:
            raise
        if printed:
            print()
        snapshot = stream.current_message_snapshot if started else None
        raise StreamFailure(exc, snapshot, finished) from exc

stream() 要搭配 with 使用,打開之後用 for 一個一個讀,讀到的就是前面那一串事件。遇到 text_delta 就把那一小段文字印出來,end="" 讓它接在同一行,flush=True 讓它馬上出現在畫面上,不要卡在緩衝區裡。

收完之後呼叫 get_final_message(),官方文件說它會把所有事件累積起來,回傳完整的 Message 物件,跟 create() 回傳的一樣,所以 run_agent() 函式後面處理回應的程式碼幾乎不用動。

finished 跟 started 這兩個變數是替中斷準備的,正常跑完的時候用不到。每收到一個 content_block_stop,finished 就記下「到這一塊為止都是完整的」。started 記的是有沒有收到 message_start。SDK 會一邊收事件、一邊把目前收到的內容拼成一份還沒完成的回應,這份東西叫做快照,用 current_message_snapshot 就能拿到。不過 SDK 要收到 message_start 之後才會開始拼,以我這次用的 SDK 1.9.0 來說,在那之前去讀 current_message_snapshot 會直接失敗。所以第一個事件還沒到就按下 Ctrl-C 的話,只能回傳 None。

官方文件的範例用的是另一種更簡單的寫法:

for text in stream.text_stream:
    print(text, end="", flush=True)

text_stream 只給文字,每次拿到的就是一個 text_delta 裡的那一小段字串。我沒用它是因為我還需要 content_block_stop 來數完整的積木,如果只是做純聊天,這樣寫就夠了。

印到一半按下 Ctrl-C?

正常跑完的時候,get_final_message() 會給我們一則完整的回應。可是如果回答印到一半就按下 Ctrl-C,回應還沒送完,get_final_message() 就拿不到了,這時候能用的只剩前面講的快照,也就是 SDK 到按下去那一刻為止拼出來的半成品。這份半成品裡到底有什麼?events.py 的第二段請模型「直接建立一個 hello.py」,每收到一段工具參數,就把快照裡的 input 印出來:

# 2. 工具參數一段一段到:直接建立一個 hello.py,內容是一支印出 hello 的 Python 程式,不用先查看目錄
​
  收到 ''
      快照裡的 input:{}
  收到 '{"file_'
      快照裡的 input:{}
  收到 'path": '
      快照裡的 input:{}
  收到 '"hello.'
      快照裡的 input:{}
  收到 'py"'
      快照裡的 input:{'file_path': 'hello.py'}
  收到 ', "cont'
      快照裡的 input:{'file_path': 'hello.py'}
  收到 'ent": "'
      快照裡的 input:{'file_path': 'hello.py'}
  收到 'print(\'
      快照裡的 input:{'file_path': 'hello.py'}
  收到 '"hello\'
      快照裡的 input:{'file_path': 'hello.py'}
  收到 '")\n"}'
      快照裡的 input:{'file_path': 'hello.py', 'content': 'print("hello")\n'}
​
  收完之後:write_file({'file_path': 'hello.py', 'content': 'print("hello")\n'})

"hello. 到的時候 input 還是空的,SDK 在解析半截的 JSON 時,會先略過還沒寫完的字串。等 py" 一到,file_path 就出現了,可是 content 還要再過五段才到。中間這段時間,快照裡的許願單看起來就是一張完整的 write_file,有檔名,只是沒有內容。

第三段是同一題再跑一次,在「有檔名、還沒有內容」的那一刻停下來,看快照裡留下什麼:

# 3. 參數送到一半就停下來:直接建立一個 hello.py,內容是一支印出 hello 的 Python 程式,不用先查看目錄
​
  收到第 5 段工具參數時停下來,完整的積木有 0 塊
    [0] tool_use:name=write_file input={'file_path': 'hello.py'}
    stop_reason=None

這就是按下 Ctrl-C 那一刻手上的東西。stop_reason 是 None,這張單子還沒寫完,但只看 input 分辨不出來。如果把這張單子放進日記,模型下一輪會看到自己開過一張「建立 hello.py,但沒有內容」的單子,它可能會以為檔案已經建好了,也可能照著這張不完整的單子再做一次,不管哪一種都不是我們要的。

官方文件在講串流中斷之後怎麼接續的時候也提到:

Tool use and extended thinking blocks cannot be partially recovered. You can resume streaming from the most recent text block.

工具積木沒辦法只救回一半,所以 stream_turn() 回報被中斷的時候 run_agent() 函式只留下已經完整的積木,finished 就是前面數出來的完整積木數量:

if interrupted:
    # 只留已經完整的積木,組到一半的工具單丟掉
    usable = list(resp.content[:finished]) if resp is not None else []
    if usable:
        history.append({"role": "assistant", "content": usable})
        seal_dangling_tool_use(history)
    message = "(這一輪被中斷了,還沒完成的部分沒有留下)"
    history.append({"role": "assistant", "content": message})
    return say(message)

按下去的時間點有三種:

  1. 文字講到一半,finished 是 0,模型這一輪的回覆一塊都不留,日記裡只記下「(這一輪被中斷了,還沒完成的部分沒有留下)」這句中斷提示。畫面上的半句話不會進日記,所以模型下一輪看到的會跟畫面上的不一樣,如果接著追問那半句話,模型可能接不上。

  2. 文字講完、工具參數還在送,finished 是 1,留下文字那一塊,半截的工具單丟掉。上面那次實跑的第一塊就是工具單,所以 finished 是 0,連文字都沒有。

  3. 工具單也送完了才按,兩塊都留下,seal_dangling_tool_use() 函式會替那張單補上「錯誤:使用者中斷,這個工具的執行狀態不明;重試前請先檢查現況。」。這是第 12 天的規矩,日記裡每一張許願單後面,都要有一則對應的結果緊跟著。

串流開始之後才收到 API 錯誤,例如前面文件提到的 overloaded_error,stream_turn() 會改丟出一個 StreamFailure 例外,把錯誤、快照跟 finished 一起帶回 run_agent() 函式,run_agent() 收到之後只留下已經完成的文字,工具單一律不執行,也不寫進日記。這一輪是出錯中斷的,就算工具單看起來已經送完,我也不想冒險去執行它,最後再把錯誤訊息記進去:

except StreamFailure as failure:
    # 串流已開始才失敗:留下完成的文字,工具單一律不執行也不寫進日記
    usable = (
        list(failure.snapshot.content[:failure.finished])
        if failure.snapshot is not None
        else []
    )
    usable = [block for block in usable if block_type(block) == "text"]
    if usable:
        history.append({"role": "assistant", "content": usable})
    message = describe_api_error(failure.cause)
    history.append({"role": "assistant", "content": message})
    return say(message)

不管是 Ctrl-C 還是 API 錯誤,畫面上都可能已經印出一些沒講完的文字,這些不會留在日記裡。我寧可少留一點,也不要讓半截的工具單混進下一輪。

events.py 那三段輸出,都是我在 2026 年 9 月 30 日實際跑的結果,三段都會真的呼叫 API,加起來不到 0.02 美元。模型每次的回答跟 JSON 切的位置都可能不一樣,你自己跑的時候,事件的數量、切法跟停下來的位置都可能跟我的不同。

顯示交給 run_agent

之前是 run_agent() 函式回傳文字,再由外面的 main() 迴圈負責印出來:

print("KeSi >", run_agent(client, history))

現在文字在串流的時候就印出來了,main() 再印一次就變兩份,可是錯誤訊息、圈數上限那些不是串流來的,還是得要印出來,所以我把顯示的工作全部收進 run_agent() 函式:

def say(message):
    """串流之外的訊息由這裡統一輸出,回傳值只給程式用。"""
    print(f"KeSi > {message}")
    return message

API 錯誤、達到圈數上限、工具連續失敗、被中斷,這些不是串流印出來的訊息全部改成 return say(...),main() 就只剩下呼叫:

# 顯示全部交給 run_agent,正常回應在串流時就已經印出來了
run_agent(client, history)

回傳值還在,只是不再拿來印,留給程式用,例如驗收腳本要檢查它回了什麼。

驗收串流!

驗收一樣不花錢,打的是第 12 天做的那台假的 API 伺服器,也就是 mock server,它看到請求裡有 "stream": true,就會把同一則訊息拆成 SSE 事件送回來。SSE 是 Server-Sent Events,伺服器把回應切成一則一則的事件持續送過來,每一則就是一行 event: 加一行 data:,再用一個空行隔開:

def event(self, name, data):
    body = f"event: {name}\ndata: {json.dumps(data)}\n\n"
    self.wfile.write(body.encode())
    self.wfile.flush()

不過這台 mock 有兩種情況做不出來。它會把整份工具參數放在同一段 input_json_delta 送出,而且它的錯誤都發生在串流開始之前。所以「參數送到一半」跟「串流途中出錯」這兩種情況,驗收改用一個自己寫的假 stream 物件,它用起來跟 SDK 的串流一樣,但跑到指定的事件就會丟出 KeyboardInterrupt 或 API 錯誤,假 stream 的快照也照前面實際看到的樣子,工具單的 input 已經有一個看起來完整的 file_path:

uv run examples/day17/check.py
# 1. 串流跑完的日記結構正確
KeSi > 我看一下這個檔案。
  [執行工具] read_file({'file_path': '不存在的檔案.txt'})
KeSi > 我看一下這個檔案。
  [執行工具] read_file({'file_path': '不存在的檔案.txt'})
KeSi > 我看一下這個檔案。
  [執行工具] read_file({'file_path': '不存在的檔案.txt'})
KeSi > 錯誤:工具連續 3 輪全部失敗,在第 3 圈停下來,請換個方式再試。
PASS 日記是 user、assistant 交錯  共 8 則
PASS 沒有懸空的許願單
​
# 2. 文字是一段一段印出來的
KeSi > (mock)好了。
PASS 文字分成多次印出  文字分 2 次印出
​
# 3. 半截的許願單不會進日記
KeSi > 我來讀一下
KeSi > (這一輪被中斷了,還沒完成的部分沒有留下)
PASS 半截的工具單沒有進日記  留下 1 塊,其中工具單 0 張
PASS 中斷後沒有懸空的許願單
KeSi > (這一輪被中斷了,還沒完成的部分沒有留下)
PASS message_start 前中斷,不會去讀不存在的快照  日記 2 則
​
# 4. 串流途中出錯,日記仍然完整
KeSi > 我來讀一下
KeSi > 錯誤:FakeStreamError:mock stream error
PASS 串流途中出錯只留下完成的文字  日記 3 則
​
# 5. 補完之後的日記能再送出一次請求
PASS 補完之後的日記送到 mock 也收到回應  補完之後 3 則,stop_reason=end_turn
​
# 8/8 通過

第 3 組試的就是剛剛講的兩種中斷。參數送到一半的時候按下去,文字留下,看起來完整的半截工具單沒有進日記。第一個事件還沒到就按下去,程式不會去讀還不存在的快照,模型的回覆一塊都沒留,日記裡只有問題跟那句中斷提示,所以是 2 則。

有比較快嗎?

我在 examples/day17/timing.py 拿兩個題目分別用串流跟非串流模式交錯著各跑 3 次,記下第一段文字什麼時候到以及整則訊息什麼時候收完。短的那題就是前面讀 note.txt 那一題,長的那題請模型用大約 600 個中文字介紹 list 跟 tuple:

uv run --env-file .env examples/day17/timing.py
# 短:幫我讀一下 note.txt 的內容
  第 1 次  串流   第一段文字  1.31s  收完  1.72s  output 77
  第 1 次  不串流  第一段文字  1.99s  收完  1.99s  output 54
  第 2 次  串流   第一段文字   (無)  收完  1.19s  output 58
  第 2 次  不串流  第一段文字  1.50s  收完  1.50s  output 76
  第 3 次  串流   第一段文字   (無)  收完  1.27s  output 58
  第 3 次  不串流  第一段文字  1.37s  收完  1.37s  output 76
​
# 長:不用查任何檔案,直接用大約 600 個中文字介紹 Python 的 list 跟 tuple 差在哪
  第 1 次  串流   第一段文字  0.96s  收完  6.62s  output 614
  第 1 次  不串流  第一段文字  5.69s  收完  5.69s  output 546
  第 2 次  串流   第一段文字  1.23s  收完  6.29s  output 565
  第 2 次  不串流  第一段文字  6.61s  收完  6.61s  output 589
  第 3 次  串流   第一段文字  1.00s  收完  5.95s  output 546
  第 3 次  不串流  第一段文字  5.49s  收完  5.49s  output 566

短的那題幾乎沒差,兩邊都在 1 到 2 秒之間收完。而且串流的第 2、3 次根本沒有文字,模型沒先講話就直接開了工具單,這種時候串流也沒有東西可以先印。

不過長的那題就差比較多了,串流大約 1 秒就看到第一段文字,不串流要等 5.5 到 6.6 秒才一次出現。整則收完的時間兩邊差不多,串流是 5.95 到 6.62 秒,不串流是 5.49 到 6.61 秒,只跑 3 次看不出誰比較快。第一段文字早點出現的好處是,可以早點發現模型理解錯了,直接按 Ctrl-C,不用等它把整段錯的講完。

output token 每次也不太一樣,短的在 54 到 77 之間、長的在 546 到 614 之間,因為每次生成的內容本來就不同。串流跟不串流回報的是同一組 usage 欄位,錢一樣是照 token 算。

這些數字是我在 2026 年 9 月 30 日實際跑的,12 次請求加起來不到 0.1 美元。時間會受網路跟服務的負載影響,你自己跑的數字可能跟我的不同。

小結

今天把第 2 天留下的 "stream": true 打開了,程式的改動不大,create() 換成 stream(),一個一個讀事件,遇到 text_delta 就把文字印出來。

比較麻煩的是怎麼處理中斷。串流讓「回應到一半」變成一個真的會遇到的狀態,而這次用的 SDK 會把送到一半的工具參數先拼成一個看起來完整的 input,只看快照分不出它送完了沒有。所以我用 content_block_stop 數出幾塊是完整的,中斷的時候只留那幾塊。以這次 600 字那題來說,串流大約 1 秒就看到第一段文字,不串流要等 5 秒以上,整則收完的時間兩邊差不多。

不過還有幾件事今天沒有處理,例如工具參數目前沒有即時顯示,[執行工具] 那行還是等整張許願單送完才印,要邊送邊顯示就得處理不完整的 JSON,還要能改寫已經印在畫面上的東西,在第 1 天的文章就曾提到這系列的重點是 agent 的腦而不是它的臉。另外中斷的組合我也還沒全部測過,例如串流剛結束、工具還沒開始跑的那一瞬間,目前還是靠 main() 接到 KeyboardInterrupt 之後呼叫 seal_dangling_tool_use() 函式補上結果。

不過明天先來算錢錢,畢竟前面好幾天都在講 token 跟成本也實際燒了幾塊美金當學費了,但 KeSi 本身到現在還沒有記帳的功能,明天幫它把每一輪花了多少算出來。

咱們下集見 :)

合作夥伴

留言討論