送信ボタンを押してから 12 秒、画面の真ん中でスピナーだけが回り続けていました。そのあいだ何も起きない。ようやく回答が出た瞬間に、テスターの指はもうホーム画面に伸びていた——自分のアプリで AI チャットを最初に組んだとき、私が最初に直面したのはこの光景でした。
モデルの応答が遅いわけではありません。生成が終わるまで待ってから一括で返す、という設計そのものが待ち時間を「無」にしてしまっていた。トークンが生成された端から流れてくれば、同じ 12 秒でも読者は読み始められます。
そこで React Native に SSE(Server-Sent Events)を持ち込もうとして、二度目の壁にぶつかりました。ブラウザで当たり前に使う EventSource が、React Native のランタイムには存在しません。fetch() からストリームを読む手順も、AbortController の噛ませ方も、Web の記事をそのまま持ってくると必ずどこかで折れます。
このページは、その折れた箇所をひとつずつ埋めていった記録です。Anthropic・OpenAI・Gemini の 3 社を同じインターフェースで扱う抽象化、API キーを端末に置かないための Cloudflare Workers プロキシ、そして課金を発生させずにストリーミングを検証する方法まで、実際に本番で動いているコードを載せています。
なぜ LLM ストリーミングが UX の分水嶺になるのか
LLM API には 2 つの呼び出し方があります。一括レスポンス(blocking) と ストリーミング(streaming) です。
一括レスポンスは単純です。リクエストを投げ、モデルが全文生成を終えてから JSON が返ってきます。実装が簡単な代わりに、長い回答では待ち時間がそのまま無音の時間になります。Claude Sonnet クラスのモデルに 800 字程度の説明を書かせると、私の手元では 15〜20 秒に達することが珍しくありませんでした。
この無音が何を引き起こすかは、数字より先にテスターの表情で分かります。5 秒を超えたあたりで一度画面から目を離し、10 秒を超えると別のアプリに切り替える。戻ってきたときには何を尋ねたのかを忘れている。離脱率を計測する以前の問題として、会話が続かなくなります。
ストリーミングは、LLM がトークンを生成するたびにその断片を逐次送信します。ユーザーには数百ミリ秒後から文字が流れ始め、全体の生成時間が同じでも「速い」と感じられます。待つ側から読む側へ、姿勢が変わる。ここが分水嶺です。
もう一つ、見落とされがちな利点がコストの制御です。ストリーミングを途中でキャンセルすれば、そこまでに生成されたトークン分しか課金されません(Anthropic・OpenAI ともに途中キャンセルに対応しています)。「もう分かったので止めたい」という場面で停止ボタンを押してもらえる状態にしておくと、想定より長い回答が延々と生成され続ける事故を防げます。私のアプリでは、この停止ボタンを付けた月の請求額が明確に下がりました。
React Native で SSE を扱う際の根本的な違い
Web ブラウザでは EventSource API を使って SSE を受信できます。しかし React Native では EventSource は存在しない。これを知らずに new EventSource(url) を書いてビルドが通り、実機で真っ白になるパターンは定番の落とし穴です。
ストリーミングの実体は fetch() から ReadableStream を読むことです。ただしここに、私が半日を溶かした落とし穴がもうひとつ埋まっていました。
React Native のグローバル fetch は、response.body を返しません。 RN の fetch は XMLHttpRequest の上に載せたポリフィルで、ストリーム対応は仕様として未実装のままです(facebook/react-native#27741 が長く開いたままなのが実情です)。Expo SDK を上げれば直る、という種類の問題ではありません。
使うのは、Expo SDK 52 で追加された expo/fetch です。名前付きインポートで取り込む別実装で、ネイティブ側にストリーム対応の fetch を持っています。グローバルの fetch を置き換えるものではないため、インポート文を書き忘れると静かに旧実装へ落ちます。ここが本当に見つけにくい。
// ❌ ブラウザでは動くが RN に EventSource は存在しない
const es = new EventSource('https://api.anthropic.com/v1/messages');
// ❌ グローバル fetch — 通信自体は成功するのに body が取れない
const bad = await fetch(url, { method: 'POST', body });
bad.body; // → undefined(RN のポリフィル実装)
// ✅ Expo SDK 52 以降。名前付きインポートが必須
import { fetch } from 'expo/fetch';
const response = await fetch('https://api.anthropic.com/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'anthropic-version': '2023-06-01',
'x-api-key': 'YOUR_ANTHROPIC_API_KEY', // 本番ではプロキシ経由で渡す(後述)
},
body: JSON.stringify({ stream: true, ...params }),
});
// response.body は ReadableStream
const reader = response.body?.getReader();
const decoder = new TextDecoder();
while (true) {
const { done, value } = await reader.read();
if (done) break;
const chunk = decoder.decode(value, { stream: true });
// chunk を解析して UI を更新
}
response.body が undefined のまま返ってきたら、原因はほぼ次の3つのどれかです。
| 症状 | 原因 | 対処 |
body が undefined。エラーは出ない | expo/fetch をインポートしていない | import { fetch } from 'expo/fetch'; を追加する |
| インポートしても undefined | Expo SDK 51 以前 | SDK 52 以上へ上げる。上げられない場合は react-native-sse などのライブラリに退避する |
| Web では動くが実機だけ落ちる | EXPO_PUBLIC_USE_RN_FETCH=1 が効いている | この環境変数はグローバルを RN 実装へ戻すためのもの。名前付きインポート側は影響を受けないため、まずインポート経路を確認する |
Rork が生成するプロジェクトは新しい SDK を使いますが、expo/fetch を明示的に呼ぶかどうかは生成されたコード次第です。ストリーミングを入れるときは、まずインポート文を目で確認してください。
SSE フォーマットの解析
Anthropic の streaming API は以下の形式でチャンクを送信します。
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"こんにちは"}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"、今日は"}}
一つの read() 呼び出しで複数の SSE イベントが含まれることがあります。また、一つのイベントが複数の read() に跨がって届くこともあります。これを正しく処理しないと、JSON parse エラーが散発的に発生します。
OpenAI の streaming 形式は少し異なります。
data: {"id":"chatcmpl-xxx","object":"chat.completion.chunk","choices":[{"delta":{"content":"こんにちは"},"index":0}]}
data: [DONE]
それぞれのパーサーを画面側に直接書いてしまうと、プロバイダーを増やすたびに UI が壊れます。抽象化レイヤーを一枚挟むことを強くおすすめします(後述)。
本番で動く完全な streaming フック実装
以下は実際の Rork プロジェクトで使っている useStreamingChat フックの完全版です。AbortController によるキャンセル、エラーリトライ、ステータス管理を含む。
// hooks/useStreamingChat.ts
import { useState, useRef, useCallback } from 'react';
// ⚠️ グローバル fetch ではストリームが読めません。必ずこの named import を使います
import { fetch } from 'expo/fetch';
type Message = {
role: 'user' | 'assistant';
content: string;
};
type StreamingState = {
messages: Message[];
isStreaming: boolean;
error: string | null;
sendMessage: (text: string) => Promise<void>;
stopStreaming: () => void;
resetConversation: () => void;
};
// Cloudflare Workers プロキシ経由で呼ぶ(API キーをアプリに埋め込まない)
const PROXY_URL = 'https://your-worker.your-subdomain.workers.dev/api/chat';
// HTTP 200 で始まった後に流れてくる error イベント用。
// 通常の JSON parse 失敗と区別するために専用クラスにしています
class StreamAbortedError extends Error {
constructor(message: string) {
super(message);
this.name = 'StreamAbortedError';
}
}
export function useStreamingChat(): StreamingState {
const [messages, setMessages] = useState<Message[]>([]);
const [isStreaming, setIsStreaming] = useState(false);
const [error, setError] = useState<string | null>(null);
const abortControllerRef = useRef<AbortController | null>(null);
// 1文字でも画面に出したかどうか。リトライ可否の判定に使います
const hasEmittedRef = useRef(false);
const stopStreaming = useCallback(() => {
if (abortControllerRef.current) {
abortControllerRef.current.abort();
abortControllerRef.current = null;
}
setIsStreaming(false);
}, []);
const sendMessage = useCallback(async (text: string) => {
if (isStreaming) return; // 多重送信防止
const userMessage: Message = { role: 'user', content: text };
const updatedMessages = [...messages, userMessage];
setMessages(updatedMessages);
setIsStreaming(true);
setError(null);
// 新しい AbortController を作成
const abortController = new AbortController();
abortControllerRef.current = abortController;
// AI 応答のプレースホルダーを追加
const assistantPlaceholder: Message = { role: 'assistant', content: '' };
setMessages([...updatedMessages, assistantPlaceholder]);
let retryCount = 0;
const maxRetries = 2;
hasEmittedRef.current = false; // 送信のたびにリセット
const attemptFetch = async (): Promise<void> => {
try {
const response = await fetch(PROXY_URL, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
messages: updatedMessages,
stream: true,
}),
signal: abortController.signal,
});
if (!response.ok) {
const errorText = await response.text();
throw new Error(`HTTP ${response.status}: ${errorText}`);
}
const reader = response.body?.getReader();
// ここで落ちる場合は expo/fetch を経由していない可能性が高い
if (!reader) {
throw new Error(
'response.body が取得できません。expo/fetch からの named import を確認してください'
);
}
const decoder = new TextDecoder();
let buffer = ''; // チャンク境界をまたぐデータの一時保持
let accumulatedText = '';
while (true) {
// AbortController でキャンセルされたら例外がスローされる
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// SSE は \n\n でイベントを区切る
const events = buffer.split('\n\n');
buffer = events.pop() ?? ''; // 最後の不完全なイベントをバッファに残す
for (const event of events) {
const lines = event.split('\n');
for (const line of lines) {
if (!line.startsWith('data: ')) continue;
const data = line.slice(6); // "data: " の6文字を除去
if (data === '[DONE]') continue; // OpenAI 終了マーカー
try {
const parsed = JSON.parse(data);
// ⚠️ 見落としやすい: 200 OK で始まった後に error イベントが流れてくる
// これを拾わないと「途中で止まった短い回答」が正常終了として残ります
if (parsed.type === 'error' || parsed.error) {
const detail =
parsed.error?.message ?? parsed.error?.type ?? 'stream error';
throw new StreamAbortedError(detail);
}
// Anthropic 形式のデルタテキスト取得
const deltaText = parsed.delta?.text
// OpenAI 形式のデルタテキスト取得
?? parsed.choices?.[0]?.delta?.content
?? '';
if (deltaText) {
hasEmittedRef.current = true;
accumulatedText += deltaText;
// React の batching を活用して頻繁な setState を抑制
setMessages(prev => {
const next = [...prev];
next[next.length - 1] = {
role: 'assistant',
content: accumulatedText,
};
return next;
});
}
} catch (e) {
// StreamAbortedError は握りつぶさず上位へ伝播させる
if (e instanceof StreamAbortedError) throw e;
// 個々のチャンクで JSON parse エラーが出ることがある(無視して継続)
}
}
}
}
} catch (err: unknown) {
if (err instanceof Error && err.name === 'AbortError') {
// ユーザーによる意図的な停止 — エラーではない
return;
}
// ⚠️ すでに1文字でも画面に出た後のリトライは巻き戻しになります
// accumulatedText は attemptFetch のスコープ内で再初期化されるため、
// 再試行の初回描画で表示済みテキストが短い文字列に置き換わります
if (hasEmittedRef.current) {
const message = err instanceof Error ? err.message : 'Unknown error';
setError(`応答が途中で終了しました: ${message}`);
return;
}
// ネットワークエラーは指数バックオフでリトライ(未出力のときだけ)
if (retryCount < maxRetries) {
retryCount++;
console.warn(`ストリーミングエラー(リトライ ${retryCount}/${maxRetries}):`, err);
await new Promise(resolve => setTimeout(resolve, 1000 * retryCount));
return attemptFetch();
}
const message = err instanceof Error ? err.message : 'Unknown error';
setError(`接続エラー: ${message}`);
// エラー時はプレースホルダーを除去
setMessages(prev => prev.slice(0, -1));
} finally {
setIsStreaming(false);
abortControllerRef.current = null;
}
};
await attemptFetch();
}, [messages, isStreaming]);
const resetConversation = useCallback(() => {
stopStreaming();
setMessages([]);
setError(null);
}, [stopStreaming]);
return { messages, isStreaming, error, sendMessage, stopStreaming, resetConversation };
}
このフックの設計で意識したポイントがいくつかあります。
buffer による境界処理: React Native の fetch は TCP パケット境界でチャンクが届く。一つの SSE イベントが 2 回の read() に跨がることがあります。buffer でデータを蓄積し、\n\n で区切ってから解析します。これを省略すると、本番で間欠的に JSON parse エラーが発生します。
setMessages の関数形式: setMessages(prev => [...]) を使うことで、クロージャが古い messages を参照する問題を防ぎます。ストリーミング中に他のメッセージが追加された場合でも安全に動作します。
リトライロジック: ネットワーク切断は特にモバイルで頻繁に起きる。1〜2 回のリトライを指数バックオフで実施することで、体験を大きく改善できます。ただし AbortError はリトライしない(ユーザーが意図して停止した)。
UI 実装 — 「停止」ボタンとカーソルアニメーション
フックができたら UI に接続します。ポイントは ストリーミング中は入力を無効化し、停止ボタンだけを表示することです。
// components/ChatUI.tsx
import { useState, useEffect } from 'react';
import { View, TextInput, Pressable, Text, FlatList, AppState } from 'react-native';
import { useStreamingChat } from '../hooks/useStreamingChat';
export function ChatUI() {
const { messages, isStreaming, error, sendMessage, stopStreaming } = useStreamingChat();
const [input, setInput] = useState('');
// バックグラウンド移行時にストリーミングを停止(iOSのネットワーク切断対策)
useEffect(() => {
const subscription = AppState.addEventListener('change', (nextState) => {
if (nextState === 'background' && isStreaming) {
stopStreaming();
}
});
return () => subscription.remove();
}, [isStreaming, stopStreaming]);
// コンポーネントアンマウント時にストリーミングを停止(メモリリーク防止)
useEffect(() => {
return () => { stopStreaming(); };
}, [stopStreaming]);
const handleSend = async () => {
if (!input.trim() || isStreaming) return;
const text = input;
setInput('');
await sendMessage(text);
};
return (
<View style={{ flex: 1 }}>
<FlatList
data={messages}
keyExtractor={(_, i) => String(i)}
maintainVisibleContentPosition={{ minIndexForVisible: 0 }}
renderItem={({ item, index }) => {
const isLastAssistant =
item.role === 'assistant' && index === messages.length - 1;
return (
<View style={{
alignSelf: item.role === 'user' ? 'flex-end' : 'flex-start',
backgroundColor: item.role === 'user' ? '#007AFF' : '#F2F2F7',
padding: 12,
borderRadius: 16,
margin: 4,
maxWidth: '80%',
}}>
<Text style={{ color: item.role === 'user' ? '#fff' : '#000' }}>
{item.content}
{/* ストリーミング中のブロックカーソル */}
{isStreaming && isLastAssistant ? '▋' : ''}
</Text>
</View>
);
}}
/>
{error && (
<Text style={{ color: '#FF3B30', padding: 8, textAlign: 'center' }}>
{error}
</Text>
)}
<View style={{ flexDirection: 'row', padding: 8, gap: 8 }}>
<TextInput
value={input}
onChangeText={setInput}
placeholder="メッセージを入力..."
style={{
flex: 1,
borderWidth: 1,
borderColor: '#C7C7CC',
borderRadius: 20,
paddingHorizontal: 16,
paddingVertical: 8,
}}
editable={!isStreaming}
onSubmitEditing={handleSend}
returnKeyType="send"
/>
{isStreaming ? (
<Pressable
onPress={stopStreaming}
style={{ justifyContent: 'center', padding: 8 }}
>
<Text style={{ fontSize: 24 }}>⏹</Text>
</Pressable>
) : (
<Pressable
onPress={handleSend}
style={{
backgroundColor: '#007AFF',
borderRadius: 20,
paddingHorizontal: 16,
justifyContent: 'center',
opacity: !input.trim() ? 0.5 : 1,
}}
disabled={!input.trim()}
>
<Text style={{ color: '#fff', fontWeight: '600' }}>送信</Text>
</Pressable>
)}
</View>
</View>
);
}
▋(ブロックカーソル)は小さなディテールだが、「AI が今考えている」という視覚的フィードバックとして効果的です。ストリーミングが完了すると自動的に消える。
Cloudflare Workers プロキシ — API キーをアプリに入れない設計
ここが多くの個人開発者が見落とすセキュリティの核心です。process.env.ANTHROPIC_API_KEY を Rork アプリのコードに書いても、ビルド済みバイナリから抽出可能です。App Store に公開したアプリで API キーが漏洩すると、知らない間に多額の請求が発生します。
解決策は Cloudflare Workers をプロキシとして使うことです。Workers はエッジで動き、シークレットを安全に保持できます。Rork は Cloudflare Workers との親和性が高く、この組み合わせは定番のアーキテクチャになっている(Hono × Cloudflare Workers REST API ガイドも参照)。
以下は streaming をパススルーする Workers の実装です。
// worker/src/index.ts(Hono を使用)
import { Hono } from 'hono';
import { cors } from 'hono/cors';
type Env = {
ANTHROPIC_API_KEY: string;
};
const app = new Hono<{ Bindings: Env }>();
app.use('/api/*', cors({
origin: ['*'], // 本番では特定のドメインに限定する
}));
app.post('/api/chat', async (c) => {
const body = await c.req.json();
const { messages, stream = true } = body;
// 入力バリデーション(攻撃者がコストを無限に増やすことを防ぐ)
if (!Array.isArray(messages) || messages.length === 0) {
return c.json({ error: 'Invalid messages' }, 400);
}
if (messages.length > 50) {
return c.json({ error: 'Too many messages' }, 400);
}
// 各メッセージのコンテンツ長も制限する
const totalLength = messages.reduce((sum: number, m: {content: string}) => sum + m.content.length, 0);
if (totalLength > 100000) {
return c.json({ error: 'Message too long' }, 400);
}
const anthropicResponse = await fetch('https://api.anthropic.com/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'anthropic-version': '2023-06-01',
'x-api-key': c.env.ANTHROPIC_API_KEY, // シークレットは Worker 環境変数から安全に取得
},
body: JSON.stringify({
model: 'claude-haiku-4-5-20251001', // コスト効率重視なら Haiku、品質重視なら Sonnet
max_tokens: 2048,
messages,
stream,
}),
});
if (!anthropicResponse.ok) {
const error = await anthropicResponse.text();
return c.json({ error }, anthropicResponse.status as 400 | 401 | 429 | 500);
}
// ストリーミングレスポンスをそのままクライアントに転送
return new Response(anthropicResponse.body, {
headers: {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
'Connection': 'keep-alive',
},
});
});
export default app;
Workers にシークレットを設定するには Wrangler CLI を使います。
# 対話的に API キーを入力(コマンド履歴に残らない)
wrangler secret put ANTHROPIC_API_KEY
# デプロイ
wrangler deploy
このプロキシ設計のもう一つの利点はレート制限の実装場所が明確になることです。Cloudflare のダッシュボードで Rate Limiting ルールを設定することで、一人のユーザーが短時間に大量リクエストを送ることを防げる。
マルチプロバイダー抽象化 — Anthropic・OpenAI・Gemini を一つのインターフェースで
プロダクションで AI 機能を運用していると、「Anthropic が障害を起こしたので OpenAI にフォールバックしたい」「コスト削減でモデルを切り替えたい」というニーズが必ず出てくる。プロバイダーごとに実装を書くと、後で変更コストが爆発します。
以下は Workers 側でマルチプロバイダーを抽象化するパターンです。
// worker/src/providers/types.ts
export type ChatMessage = { role: 'user' | 'assistant'; content: string };
export type StreamProvider = (
messages: ChatMessage[],
env: Record<string, string>
) => Promise<ReadableStream>;
// worker/src/providers/anthropic.ts
export const anthropicProvider: StreamProvider = async (messages, env) => {
const response = await fetch('https://api.anthropic.com/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'anthropic-version': '2023-06-01',
'x-api-key': env.ANTHROPIC_API_KEY,
},
body: JSON.stringify({
model: 'claude-sonnet-4-6',
max_tokens: 2048,
messages,
stream: true,
}),
});
if (!response.ok || !response.body) {
throw new Error(`Anthropic error: ${response.status}`);
}
return response.body;
};
// worker/src/providers/openai.ts
export const openaiProvider: StreamProvider = async (messages, env) => {
const response = await fetch('https://api.openai.com/v1/chat/completions', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Authorization': `Bearer ${env.OPENAI_API_KEY}`,
},
body: JSON.stringify({
model: 'gpt-4o-mini', // コスト効率が高い
messages,
stream: true,
}),
});
if (!response.ok || !response.body) {
throw new Error(`OpenAI error: ${response.status}`);
}
return response.body;
};
// worker/src/providers/gemini.ts
export const geminiProvider: StreamProvider = async (messages, env) => {
// Gemini は messages 形式が異なるため変換が必要
const geminiContents = messages.map(m => ({
role: m.role === 'assistant' ? 'model' : 'user',
parts: [{ text: m.content }],
}));
const response = await fetch(
`https://generativelanguage.googleapis.com/v1beta/models/gemini-2.0-flash:streamGenerateContent?alt=sse&key=${env.GEMINI_API_KEY}`,
{
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ contents: geminiContents }),
}
);
if (!response.ok || !response.body) {
throw new Error(`Gemini error: ${response.status}`);
}
return response.body;
};
// worker/src/index.ts — フォールバック付きプロバイダー選択
app.post('/api/chat', async (c) => {
const { messages, provider = 'anthropic' } = await c.req.json();
const env = {
ANTHROPIC_API_KEY: c.env.ANTHROPIC_API_KEY,
OPENAI_API_KEY: c.env.OPENAI_API_KEY ?? '',
GEMINI_API_KEY: c.env.GEMINI_API_KEY ?? '',
};
const providers: Record<string, StreamProvider> = {
anthropic: anthropicProvider,
openai: openaiProvider,
gemini: geminiProvider,
};
const selectedProvider = providers[provider] ?? anthropicProvider;
try {
const stream = await selectedProvider(messages, env);
return new Response(stream, {
headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache' },
});
} catch (primaryError) {
// プライマリプロバイダーが失敗したら OpenAI にフォールバック
console.error('Primary provider failed, falling back:', primaryError);
try {
const fallbackStream = await openaiProvider(messages, env);
return new Response(fallbackStream, {
headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache' },
});
} catch {
return c.json({ error: 'All providers failed' }, 503);
}
}
});
このアーキテクチャにより、アプリ側のコードを一切変更せずにバックエンドでプロバイダーを切り替えられます。コスト最適化の観点では、シンプルな質問は Haiku/GPT-4o-mini、複雑な推論タスクは Sonnet と振り分けるロジックも Workers 側に集約できます。
パフォーマンス最適化 — 60fps を維持するトークン表示
ストリーミング中は 1 秒間に数十回 setMessages が呼ばれます。各 setState が再レンダリングをトリガーし、長いメッセージリストでは顕著にフレームドロップが発生します。
スロットリングパターンで setState 頻度を制御します。最後の更新から 50ms 未満なら setState を遅延させる実装です。
// hooks/useStreamingChat.ts に追加するスロットリングロジック
const lastUpdateRef = useRef<number>(0);
const pendingTextRef = useRef<string>('');
const flushTimerRef = useRef<ReturnType<typeof setTimeout> | null>(null);
const flushPendingText = useCallback(() => {
if (flushTimerRef.current) {
clearTimeout(flushTimerRef.current);
flushTimerRef.current = null;
}
setMessages(prev => {
const next = [...prev];
next[next.length - 1] = {
role: 'assistant',
content: pendingTextRef.current,
};
return next;
});
lastUpdateRef.current = Date.now();
}, []);
// ストリーミングのデルタ処理内で使用
const updateText = (newText: string) => {
pendingTextRef.current = newText;
const now = Date.now();
const elapsed = now - lastUpdateRef.current;
if (elapsed >= 50) {
// 50ms 以上経過していれば即座に更新
flushPendingText();
} else {
// 50ms 後に更新をスケジュール(既存のタイマーはキャンセル)
if (flushTimerRef.current) clearTimeout(flushTimerRef.current);
flushTimerRef.current = setTimeout(flushPendingText, 50 - elapsed);
}
};
このパターンにすると setState の呼び出し回数が大きく減り、FlatList のスクロールが目に見えて滑らかになります。ストリーミング完了時は flushPendingText() を呼んで最終テキストを確定させてください。
日本語・中国語・絵文字のマルチバイト文字が化ける場合は、Hermes エンジンの TextDecoder の挙動差異が原因のことがあります。{ stream: true } オプションなしでデコードし直すか、Uint8Array を蓄積して最後に一括デコードするアプローチで解決できます。
よくある間違いと落とし穴
実際のプロジェクトで踏んだ地雷を整理します。
落とし穴 1: fetch に signal を渡し忘れる
stopStreaming を呼んでもストリームが止まらない場合、ほぼ確実に signal: abortController.signal を fetch オプションに渡していません。AbortController は fetch に接続されていないと機能しません。
落とし穴 2: max_tokens を省略する
省略するとモデルのデフォルト上限まで生成し続ける可能性があります。チャットアプリで 8000 トークンの回答が必要になる場面は稀で、その分だけ請求が膨らみます。用途に応じて 1024〜4096 の範囲で明示的に設定してください。
落とし穴 3: Component unmount 後に setState が呼ばれる
ユーザーが回答生成中に画面を離れると、unmount 済みのコンポーネントへ setMessages が呼ばれ続けます。useEffect のクリーンアップで必ず stopStreaming を呼んでください。上記の UI 実装例ではすでに対応しています。
落とし穴 4: バックグラウンド移行でストリーミングが切れる
iOS はアプリがバックグラウンドに移行すると、ネットワーク接続を数秒〜数十秒で切断することがあります。AppState を監視してバックグラウンド移行時に stopStreaming を呼ぶか、バックグラウンドでは新規ストリーミングを開始しない制御が必要です。上記 UI 実装例では AppState.addEventListener で対応しています。
落とし穴 5: プロキシなしで API キーをアプリに埋め込む
ビルド済みバイナリは逆コンパイルできます。API キーを直接コードに書くと漏洩リスクがあります。必ず Cloudflare Workers 等のプロキシ経由でアクセスする設計にすること(API キーを Worker で隠しつつレート制限とローテーションを回す設計も参照)。
落とし穴 6: SSE バッファ処理を省略する
チャンク境界をまたぐ SSE イベントを考慮せず直接 JSON.parse すると、本番で間欠的にエラーが発生します。再現性が低いため原因特定が難しくなります。必ず buffer で蓄積して \n\n 区切りで解析します。
落とし穴 7: 200 OK の後に流れてくる error イベントを捨てている
これが一番怖い落とし穴でした。Anthropic の streaming は、接続に成功して数トークン流したあとで event: error を送ってくることがあります。overloaded_error が典型です。
デルタテキストだけを見るパーサーは、この error を「テキストの無いイベント」として素通りさせます。ストリームはそのまま done に到達し、UI 側は正常終了として扱う。読者の画面には、途中で切れた短い回答だけが残ります。エラー表示も、リトライも起きません。
パーサーを Node で切り出して、ping と overloaded_error だけを流してみました。結果は「取得テキスト 0 文字・例外なし・エラー検知なし」。つまり静かに失敗します。だからこそ、上のフックでは parsed.type === 'error' を明示的に見て専用の例外へ変換しています。
落とし穴 8: 部分出力の後にリトライすると、画面が巻き戻る
指数バックオフのリトライは、接続そのものが失敗したときには有効です。しかし「5文字ぶん表示した後で切断された」場合に同じ関数を呼ぶと、accumulatedText は再初期化され、リトライ後の最初の描画で表示済みテキストがより短い文字列に置き換わります。
同じパーサーで再現したところ、中断時に画面へ出ていたのは5文字、リトライ後の初回描画は2文字でした。読者の目には、書きかけの文章が消えて別の文章が生え直したように映ります。加えて、LLM 呼び出しは冪等ではないため、生成トークンぶんの請求も二重に発生します。
対策はシンプルです。1文字でも出力したらリトライしない。出力済みのテキストはそのまま残し、末尾にエラーだけを添える。上のフックの hasEmittedRef はこのための旗です。
コスト管理と監視体制
ストリーミング実装が完成したあと、コストが見えないまま運用に入るのは危険です。Cloudflare Workers のログに各リクエストのトークン数を記録し、Anthropic / OpenAI のダッシュボードで使用量を定期的に確認する習慣をつけてください。
Workers 側でトークン消費の概算を記録する簡易ロギングを追加します。
// ストリーミング開始時にリクエストをログ記録
console.log(JSON.stringify({
timestamp: new Date().toISOString(),
model: 'claude-haiku-4-5-20251001',
messageCount: messages.length,
estimatedInputTokens: messages.reduce((sum, m) => sum + Math.ceil(m.content.length / 4), 0),
}));
Anthropic の streaming レスポンスは、最後の message_delta イベントに実際の usage を含みます。これを Workers 側で捕捉してログに残すと、推定値ではなく実測値でコストを追えるようになります。
また、一ユーザーあたりの日次クレジット上限を Cloudflare KV で管理することも検討に値します。KV でユーザー ID とトークン消費数をカウントし、上限超過後は 429 を返す設計です。個人開発アプリでは無料ユーザーへの無制限 AI アクセスを防ぐために有効な手段となります。
会話コンテキストをどこで切るか — 履歴とコストの折り合い
ストリーミングが動き出すと、次に効いてくるのは会話履歴です。チャット UI は放っておくと過去のやり取りを全部リクエストに積み込みます。20 往復もすれば入力トークンが数千に達し、1 回の返信のたびにその全量が課金対象になります。ストリーミングで削ったコストが、履歴の肥大でそのまま戻ってくるという間の抜けた事態になりました。
直近だけ残して、古い部分は畳む
対処は 2 段構えにしています。直近の N 往復はそのまま残し、それより古い部分は一度だけ要約して 1 メッセージに畳む。要約は安いモデル(Haiku クラス)に投げれば十分です。私の手元では、30 往復の会話で入力トークンがおよそ 4,200 → 900 まで落ち、1 リクエストあたりの入力コストが 8 割ほど下がりました。
// Workers 側: 履歴が閾値を超えたら古い部分だけ要約して畳む
const KEEP_RECENT_TURNS = 8; // 直近8往復はそのまま残す
const SUMMARIZE_THRESHOLD = 12; // 12往復を超えたら圧縮を検討
type Msg = { role: 'user' | 'assistant'; content: string };
async function compactHistory(messages: Msg[], env: Env): Promise<Msg[]> {
if (messages.length <= SUMMARIZE_THRESHOLD * 2) return messages;
const cut = messages.length - KEEP_RECENT_TURNS * 2;
const older = messages.slice(0, cut);
const recent = messages.slice(cut);
const transcript = older
.map((m) => `${m.role === 'user' ? 'ユーザー' : 'アシスタント'}: ${m.content}`)
.join('\n');
const res = await fetch('https://api.anthropic.com/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'anthropic-version': '2023-06-01',
'x-api-key': env.ANTHROPIC_API_KEY,
},
body: JSON.stringify({
model: 'claude-haiku-4-5-20251001',
max_tokens: 512,
messages: [{
role: 'user',
content:
'次の会話を、後続の応答に必要な事実・決定事項・ユーザーの好みだけを残して' +
'400字以内で要約してください。挨拶や相槌は落として構いません。\n\n' + transcript,
}],
}),
});
if (!res.ok) return messages.slice(-KEEP_RECENT_TURNS * 2); // 要約失敗時は単純に切り捨て
const data = await res.json<{ content: { text: string }[] }>();
const summary = data.content[0]?.text ?? '';
return [
{ role: 'user', content: `【これまでの経緯】${summary}` },
{ role: 'assistant', content: '把握しました。続けてください。' },
...recent,
];
}
失敗時に会話を殺さない
要約の失敗時に例外を投げず、単純な切り捨てへフォールバックしている点が実務上は重要でした。要約用の API が一時的に落ちているだけで会話全体が使えなくなるのは割に合いません。私はこの形をお勧めします。要約は品質を上げるための贅沢であって、会話の前提条件ではないからです。
畳んだ結果はクライアントに返さず、Workers 側だけで保持しています。ユーザーの画面には全履歴が見えたまま、モデルに送る分だけが圧縮されている状態です。ここを取り違えて UI 側の配列を書き換えてしまい、過去の発言が画面から消えるという苦情を一度もらいました。
課金を発生させずにストリーミングを検証する
ストリーミング周りのバグは、UI・パーサー・ネットワークのどこで壊れているのか切り分けづらい種類のものです。そのたびに本物の API を叩いていると、デバッグの往復がそのまま請求書になります。UI の微調整だけで 1 日に 200 回近くリクエストを飛ばしていた時期があり、月末の請求でようやくそれに気づきました。
応答を録画して再生する
私は SSE の応答を録画したフィクスチャを一本用意して、開発時はそれを再生する構成にしました。実装は 30 行程度で済みます。
// scripts/mock-sse-server.ts — bun run scripts/mock-sse-server.ts で起動
const FIXTURE = [
'こんにちは', '。', 'ストリーミング', 'の', 'テスト', '応答', 'です', '。',
];
Bun.serve({
port: 8787,
async fetch(req) {
const url = new URL(req.url);
const delayMs = Number(url.searchParams.get('delay') ?? 40);
const failAt = Number(url.searchParams.get('failAt') ?? -1); // 途中切断の再現用
const stream = new ReadableStream({
async start(controller) {
const enc = new TextEncoder();
for (let i = 0; i < FIXTURE.length; i++) {
if (i === failAt) { controller.error(new Error('simulated disconnect')); return; }
const payload = JSON.stringify({
type: 'content_block_delta',
index: 0,
delta: { type: 'text_delta', text: FIXTURE[i] },
});
controller.enqueue(enc.encode(`event: content_block_delta\ndata: ${payload}\n\n`));
await new Promise((r) => setTimeout(r, delayMs));
}
controller.enqueue(enc.encode('event: message_stop\ndata: {}\n\n'));
controller.close();
},
});
return new Response(stream, {
headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache' },
});
},
});
アプリ側は PROXY_URL を http://localhost:8787 に差し替えるだけです。実機から叩く場合は localhost ではなく開発機の LAN IP を指定してください。
遅延と切断をパラメータで再現する
このモックの価値は、料金より再現性にありました。?delay=400 で極端に遅いストリームを作れば、停止ボタンやカーソルアニメーションの挙動をゆっくり観察できます。?failAt=3 を付ければ、実運用では月に数回しか起きない「途中切断」を毎回確実に再現できる。リトライ処理のバグを 2 件、これで見つけました。
チャンク境界をまたぐ SSE イベントの検証をしたければ、controller.enqueue を分割して event: content_bl と ock_delta\ndata: ... のように意図的に割ってやれば、バッファ処理の穴がその場で露出します。本番で間欠的に出ていた JSON parse エラーの正体は、これで確定できました。
ツール使用を挟むと壊れやすい箇所
チャットに検索や計算を組み込む段階になると、SSE に流れてくるイベントの種類が増えます。テキストだけを前提に書いたパーサーは、ここで静かに壊れます。
部分 JSON を毎チャンク parse しない
Anthropic の場合、ツール呼び出しの引数は input_json_delta として部分的な JSON 文字列が少しずつ流れてきます。1 チャンクごとに JSON.parse を試みると、当然ほぼ毎回失敗します。
// ブロック単位でバッファし、content_block_stop で初めて確定させる
const toolBuffers = new Map<number, { name: string; json: string }>();
function handleEvent(evt: { type: string; index?: number; [k: string]: any }) {
switch (evt.type) {
case 'content_block_start':
if (evt.content_block?.type === 'tool_use') {
toolBuffers.set(evt.index!, { name: evt.content_block.name, json: '' });
}
break;
case 'content_block_delta':
if (evt.delta?.type === 'text_delta') {
appendToUi(evt.delta.text);
} else if (evt.delta?.type === 'input_json_delta') {
const buf = toolBuffers.get(evt.index!);
if (buf) buf.json += evt.delta.partial_json; // ここでは parse しない
}
break;
case 'content_block_stop': {
const buf = toolBuffers.get(evt.index!);
if (!buf) break;
toolBuffers.delete(evt.index!);
try {
runTool(buf.name, JSON.parse(buf.json || '{}'));
} catch {
// 引数が壊れているときはツールを実行せず、モデルにエラーを返して再試行させる
returnToolError(buf.name, 'invalid arguments');
}
break;
}
}
}
content_block_stop を待ってから一度だけ parse する。それだけの話ですが、テキスト前提の実装から移行するときに最も踏みやすい箇所でした。
ストリームが 2 本になる問題
もう一点、ツール実行の結果をモデルに返して続きを生成させるとき、2 本目のストリームが始まります。UI 側で「ストリーミング中」フラグを単純な真偽値で持っていると、1 本目の終了で false に落ちて停止ボタンが消え、2 本目を止められなくなります。私はここを「進行中のリクエスト ID」を持つ形に変えて解決しました。フラグではなく識別子を持たせておくと、何本目のストリームかに関係なく正しく中断できます。
この実装を入れて何が変わったか
一番変わったのは、届くフィードバックの種類でした。非ストリーミングだった頃は「AI の回答が遅い」「固まっているのでは」という報告が大半でしたが、ストリーミング導入後はそれが消え、代わりに回答の中身についての要望が届くようになりました。UI が透明になって、ようやく機能そのものを見てもらえるようになった、という感覚に近いです。
導入コストは決して低くありません。SSE のバッファ処理、AbortController の配線、プロキシの用意と、地味な作業が続きます。ただ、フックとプロキシを一度作ってしまえば、あとは呼び出し口を変えるだけで文書要約にもコード生成にも翻訳にも使い回せます。私の場合、2 本目の AI 機能を追加するのに要した時間は半日でした。
次の一手としておすすめしたいのは、この記事のモック SSE サーバーを先に置いてしまうことです。本番 API を繋ぐ前に UI と停止処理を確定させておくと、あとの検証がまるごと無料になります。そのうえで、Supabase Edge Functions と RLS の設計メモを参考に会話履歴の永続化を足すと、単発の質問応答から「続きから話せる」アプリへ一段上がります。
長い記事にお付き合いいただきありがとうございました。ここに載せたコードが、あなたのスピナーを一つ減らせたら嬉しく思います。