Bedrock Agentsで本番が崩壊した話|マルチステップ設計の地雷と対策

Bedrock Agentsの複数ステップワークフローで本番環境が落ちた経験から、状態管理・メモリ爆発・無限ループの実装パターンと対策を解説します。

はじめに:Bedrock Agentsで本番が落ちた話

先日、プロジェクトでBedrock Agentsを使ったマルチステップのワークフロー自動化に取り組んでいたんですが、本番環境で思わぬトラブルに直面しました。複数のステップを連鎖させるAIエージェント構成で、メモリ使用量が爆発的に増えたり、エージェントが無限ループに陥ったり、ステップ間の状態受け渡しが失敗したり……。当初は「Bedrock Agentsなら簡単でしょ」と甘く見ていたんですが、実際に本番環境で複数ステップを組み合わせると、複雑さが指数関数的に増えてくるんですよね。

8ヶ月の本番運用を通じて、僕たちのチームが学んだマルチステップ設計の実装パターンと落とし穴をここに書き留めておきます。これから同じ轍を踏む人の参考になれば幸いです。

Bedrock Agentsのマルチステップ構成図

flowchart TB
    subgraph VPC["AWS VPC"]
        subgraph Bedrock["Bedrock Service"]
            Agent1["Agent Step 1<br/>Planning"]
            Agent2["Agent Step 2<br/>Execute"]
            Agent3["Agent Step 3<br/>Validate"]
        end
        
        subgraph Memory["State Management"]
            SessionMemory["Session Memory<br/>Context Buffer"]
            DynamoDB[("DynamoDB<br/>State Store")]
        end
        
        subgraph Tools["Tool Integrations"]
            Lambda1["Lambda Tool 1"]
            Lambda2["Lambda Tool 2"]
            S3[("S3<br/>Artifacts")]
        end
        
        subgraph Monitoring["Observability"]
            CloudWatch["CloudWatch Logs"]
            XRay["X-Ray Tracing"]
        end
    end
    
    Client["Client App"]
    Client -->|Invoke Session| Agent1
    Agent1 -->|Next Step| Agent2
    Agent2 -->|Next Step| Agent3
    Agent3 -->|Return Result| Client
    
    Agent1 -->|Store Context| SessionMemory
    Agent2 -->|Update State| DynamoDB
    Agent3 -->|Read State| DynamoDB
    
    Agent2 -->|Call Tools| Lambda1
    Agent2 -->|Call Tools| Lambda2
    Lambda1 -->|Write| S3
    Lambda2 -->|Write| S3
    
    Agent1 --> CloudWatch
    Agent2 --> XRay
    Agent3 --> CloudWatch

実装の流れとしては、Agent Step 1でタスク計画を立てて、Step 2で実際の処理(複数のLambda Toolを呼び出し)を実行し、Step 3で結果検証をする——といった具合です。ここまで聞くと単純に見えますが、本番ではこの単純な構成がすぐに複雑化します。

マルチステップで本番が落ちた原因:メモリとコンテキスト喪失

1. エージェントのメモリ爆発

Bedrock Agentsを複数ステップで連鎖させると、各ステップで過去のやり取りが全部メモリに蓄積されていくんです。初期段階では「まあ数MB程度だろう」と思ってたんですが、実際には1セッションあたりのコンテキストが100MB近くまで膨れ上がることもありました。

僕たちの場合、ステップ1で大量のテキストデータを処理して、ステップ2でそのすべてのやり取りを参照しながら次のステップに進む——という構成を取ってました。結果として、Lambdaのメモリリミットに引っかかるか、処理時間が指数関数的に増加していったわけです。

対策としては、各ステップ完了後に不要なコンテキストを明示的に削除するというアプローチを取りました。以下は実装例です:

import json
from boto3 import client

bedrockagent = client('bedrock-agent-runtime')
dynamodb = client('dynamodb')


class MultiStepAgent:
    def __init__(self, session_id: str, dynamodb_table: str):
        self.session_id = session_id
        self.table_name = dynamodb_table
        self.max_context_size = 50000  # 50KB制限
        
    def invoke_step(self, agent_id: str, step_name: str, input_text: str):
        """マルチステップを実行し、メモリを管理"""
        try:
            # 現在のセッション状態を取得
            session_state = self._get_session_state()
            
            # ステップ実行
            response = bedrockagent.invoke_agent(
                agentId=agent_id,
                sessionId=self.session_id,
                inputText=input_text,
                sessionState={
                    'invocationId': session_state.get('invocationId'),
                    'returnControlInvocationResults': session_state.get('results', [])
                }
            )
            
            # ステップ結果を処理
            step_result = self._process_response(response)
            
            # コンテキストサイズチェック
            if self._estimate_context_size(session_state) > self.max_context_size:
                # 古いコンテキストをアーカイブして削除
                self._archive_and_prune_context()
            
            # 完了したステップの詳細情報をS3にアーカイブ
            self._archive_step_details(step_name, step_result)
            
            # DynamoDBに最小限の状態のみ保持
            self._update_minimal_state(step_name, step_result)
            
            return step_result
            
        except Exception as e:
            print(f"Step {step_name} failed: {str(e)}")
            # 失敗時のロールバック処理
            self._rollback_step(step_name)
            raise
    
    def _estimate_context_size(self, state: dict) -> int:
        """コンテキストサイズを推定"""
        return len(json.dumps(state).encode('utf-8'))
    
    def _archive_and_prune_context(self):
        """古いコンテキストを削除"""
        state = self._get_session_state()
        
        # 不要なキーを削除(詳細なやり取りログなど)
        if 'detailed_interactions' in state:
            del state['detailed_interactions']
        
        if 'intermediate_results' in state:
            # 最新のものだけ保持
            state['intermediate_results'] = state['intermediate_results'][-1:]
        
        self._update_session_state(state)
    
    def _get_session_state(self) -> dict:
        """DynamoDBからセッション状態を取得"""
        response = dynamodb.get_item(
            TableName=self.table_name,
            Key={'sessionId': {'S': self.session_id}}
        )
        return response.get('Item', {})
    
    def _update_session_state(self, state: dict):
        """セッション状態を更新(最小限の情報のみ)"""
        dynamodb.put_item(
            TableName=self.table_name,
            Item={'sessionId': {'S': self.session_id}, 'state': {'S': json.dumps(state)}}
        )
    
    def _archive_step_details(self, step_name: str, result: dict):
        """ステップの詳細情報をS3にアーカイブ"""
        s3 = client('s3')
        archive_key = f"agent-sessions/{self.session_id}/{step_name}/result.json"
        s3.put_object(
            Bucket='bedrock-agent-archives',
            Key=archive_key,
            Body=json.dumps(result)
        )
    
    def _update_minimal_state(self, step_name: str, result: dict):
        """最小限の状態のみDynamoDBに保持"""
        minimal_state = {
            'completed_steps': step_name,
            'last_result_summary': str(result)[:500],  # 最初の500文字だけ
            'timestamp': str(__import__('datetime').datetime.now())
        }
        self._update_session_state(minimal_state)
    
    def _process_response(self, response) -> dict:
        """Bedrock Agentsの応答を処理"""
        events = response.get('completion', [])
        result = {}
        
        for event in events:
            if 'actionGroupInvocation' in event:
                action = event['actionGroupInvocation']
                result['action'] = action
            elif 'finalResponse' in event:
                result['final'] = event['finalResponse']
        
        return result
    
    def _rollback_step(self, step_name: str):
        """失敗したステップのロールバック"""
        # 失敗したステップをマーク
        state = self._get_session_state()
        state['failed_step'] = step_name
        self._update_session_state(state)

2. エージェント無限ループ問題

Bedrock Agentsにツール呼び出しを許可していると、エージェント自体が「次は何をするか」を決定し続けるんです。で、その判断が微妙に間違ってると、同じツールを何度も呼び出し続けるんですよ。僕たちのテストでは、1セッションで同じLambda関数が200回以上呼ばれたケースもありました。

対策はステップごとの呼び出し回数制限明示的な終了条件を設定することです。正直、これ設定しないと本番で確実に落ちます:

class StepWithGuards:
    def __init__(self, max_tool_calls: int = 5, timeout_seconds: int = 60):
        self.max_tool_calls = max_tool_calls
        self.timeout_seconds = timeout_seconds
        self.tool_call_count = 0
        self.start_time = None
    
    def invoke_with_guards(self, agent_id: str, session_id: str, input_text: str) -> dict:
        """ガード付きでステップを実行"""
        import time
        from datetime import datetime, timedelta
        
        self.start_time = datetime.now()
        timeout = self.start_time + timedelta(seconds=self.timeout_seconds)
        self.tool_call_count = 0
        
        client = __import__('boto3').client('bedrock-agent-runtime')
        cloudwatch = __import__('boto3').client('cloudwatch')
        
        try:
            while datetime.now() < timeout and self.tool_call_count < self.max_tool_calls:
                # LLMの思考ループを実行
                response = client.invoke_agent(
                    agentId=agent_id,
                    sessionId=session_id,
                    inputText=input_text
                )
                
                # ツール呼び出しカウント
                for event in response.get('completion', []):
                    if 'actionGroupInvocation' in event:
                        self.tool_call_count += 1
                        
                        # 上限に達したら強制終了
                        if self.tool_call_count >= self.max_tool_calls:
                            print(f"Max tool calls ({self.max_tool_calls}) reached")
                            return {'status': 'max_calls_reached'}
                
                # 最終応答で終了
                if any('finalResponse' in event for event in response.get('completion', [])):
                    break
                
                # タイムアウトチェック
                if datetime.now() > timeout:
                    print(f"Step timeout after {self.timeout_seconds}s")
                    return {'status': 'timeout'}
            
            # メトリクス記録
            cloudwatch.put_metric_data(
                Namespace='BedrockAgents',
                MetricData=[
                    {
                        'MetricName': 'ToolCallsPerStep',
                        'Value': self.tool_call_count,
                        'Unit': 'Count'
                    }
                ]
            )
            
            return response
            
        except Exception as e:
            print(f"Error in step: {str(e)}")
            raise

ステップ間の状態管理:失敗からの学び

前ステップの結果を次のステップで正しく引き継ぐ

最初のうち、僕たちは「各ステップの結果をSession Stateに入れときゃいいだろう」と思ってました。でも実際には、LLMのコンテキストウインドウに収まらないサイズの結果が出てくることがあるんです。特にCSVやJSONの大量データを処理してる場合、その結果をそのままコンテキストに入れるとトークン数がオーバーフローしちゃいます。

データサイズ別のハンドリング方法は以下のようになります:

データサイズ保存先用途メリット
〜1KBDynamoDBコンテキストに含める高速アクセス
1KB〜5MBS3 + 要約版をDB要約をコンテキストに含めるバランス型
5MB〜S3のみ参照URLのみコンテキストに含めるメモリ効率的

対策の実装は以下の通りです:

class StateTransitionManager:
    def __init__(self, s3_bucket: str, dynamodb_table: str):
        self.s3 = __import__('boto3').client('s3')
        self.dynamodb = __import__('boto3').client('dynamodb')
        self.s3_bucket = s3_bucket
        self.dynamodb_table = dynamodb_table
    
    def pass_result_to_next_step(self, session_id: str, from_step: str, to_step: str, result: dict) -> str:
        """
        前ステップの結果を次のステップに渡す
        大きいデータはS3に保存し、参照URLをContextに含める
        """
        result_size = len(__import__('json').dumps(result).encode('utf-8'))
        
        # 5MB以上はS3に保存
        if result_size > 5 * 1024 * 1024:
            s3_key = f"{session_id}/{from_step}_to_{to_step}/result.json"
            self.s3.put_object(
                Bucket=self.s3_bucket,
                Key=s3_key,
                Body=__import__('json').dumps(result),
                ContentType='application/json'
            )
            
            # 次のステップ用に要約版を作成
            summary = self._create_summary(result)
            
            # DynamoDBに参照情報を保存
            self.dynamodb.put_item(
                TableName=self.dynamodb_table,
                Item={
                    'sessionId': {'S': session_id},
                    'stepTransition': {'S': f"{from_step}-{to_step}"},
                    's3Location': {'S': f"s3://{self.s3_bucket}/{s3_key}"},
                    'summary': {'S': __import__('json').dumps(summary)},
                    'originalSize': {'N': str(result_size)}
                }
            )
            
            # 次のステップに渡すプロンプト
            return self._create_context_instruction(summary, s3_key)
        else:
            # 小さいデータはそのまま保存
            self.dynamodb.put_item(
                TableName=self.dynamodb_table,
                Item={
                    'sessionId': {'S': session_id},
                    'stepTransition': {'S': f"{from_step}-{to_step}"},
                    'fullResult': {'S': __import__('json').dumps(result)}
                }
            )
            return __import__('json').dumps(result)
    
    def _create_summary(self, result: dict) -> dict:
        """大きい結果から重要な情報だけ抽出。地味に手作業になる部分"""
        summary = {}
        
        if 'records' in result:
            # リスト型の場合は最初と最後のレコード + 統計情報
            records = result['records']
            summary['totalCount'] = len(records)
            summary['firstRecord'] = records[0] if records else None
            summary['lastRecord'] = records[-1] if len(records) > 1 else None
        elif isinstance(result, dict):
            # 各キーの値のタイプと最初の部分だけを記録
            for k, v in result.items():
                if isinstance(v, str) and len(v) > 1000:
                    summary[k] = v[:1000] + '...'
                elif isinstance(v, (list, dict)):
                    summary[k] = f"{type(v).__name__} with {len(v)} items"
                else:
                    summary[k] = v
        
        return summary
    
    def _create_context_instruction(self, summary: dict, s3_key: str) -> str:
        """次のステップに渡すコンテキスト指示"""
        return f"""
前のステップの結果が大きいため、要約版と参照URLを提供します。

【要約】
{__import__('json').dumps(summary, ensure_ascii=False, indent=2)}

【詳細データ】
以下のS3パスに詳細なデータが保存されています:
s3://{self.s3_bucket}/{s3_key}

必要に応じてこのS3パスを参照してください。
"""

エラーハンドリングと復旧戦略

マルチステップだからこそ、途中で失敗する確率が高まります。ステップ1は成功したけどステップ2で失敗した場合、どうやって復旧するかが重要なんです。個人的には、ここの実装に時間をかけることで本番安定性が劇的に上がる経験をしました。

class MultiStepRecovery:
    def __init__(self, dynamodb_table: str, sns_topic_arn: str):
        self.dynamodb = __import__('boto3').client('dynamodb')
        self.sns = __import__('boto3').client('sns')
        self.table = dynamodb_table
        self.topic = sns_topic_arn
    
    def execute_with_recovery(self, session_id: str, steps: list) -> dict:
        """
        複数ステップを実行し、失敗時の自動復旧を試みる
        """
        completed_steps = []
        
        for i, step in enumerate(steps):
            step_name = step['name']
            
            try:
                print(f"Executing step {i+1}/{len(steps)}: {step_name}")
                result = self._execute_single_step(session_id, step)
                completed_steps.append({
                    'step': step_name,
                    'status': 'success',
                    'result': result
                })
                
            except Exception as e:
                print(f"Step {step_name} failed: {str(e)}")
                
                # 復旧を試みる
                recovery_result = self._attempt_recovery(session_id, step_name, e, completed_steps)
                
                if recovery_result['recovered']:
                    completed_steps.append({
                        'step': step_name,
                        'status': 'recovered',
                        'result': recovery_result['result']
                    })
                else:
                    # 復旧できなかった場合はアラート
                    self._send_alert(session_id, step_name, e, completed_steps)
                    return {
                        'status': 'failed',
                        'failedStep': step_name,
                        'completedSteps': completed_steps
                    }
        
        return {
            'status': 'success',
            'completedSteps': completed_steps
        }
    
    def _attempt_recovery(self, session_id: str, failed_step: str, error: Exception, completed_steps: list) -> dict:
        """失敗したステップの復旧を試みる"""
        error_msg = str(error)
        
        # エラーの種類別に復旧戦略を選択
        if 'timeout' in error_msg.lower():
            # タイムアウトの場合は部分結果で進める
            print(f"Timeout detected, attempting to proceed with partial results")
            return {'recovered': True, 'result': self._get_partial_result(session_id, failed_step)}
        
        elif 'rate_limit' in error_msg.lower():
            # レート制限の場合は待機して再試行
            print(f"Rate limit hit, retrying after delay")
            __import__('time').sleep(5)
            # 再試行ロジック
            return {'recovered': True, 'result': {}}
        
        elif 'memory' in error_msg.lower() or 'context' in error_msg.lower():
            # メモリ不足の場合は入力を削減して再試行
            print(f"Memory issue, reducing context and retrying")
            return {'recovered': True, 'result': {}}
        
        else:
            # その他のエラーは復旧不可
            return {'recovered': False}
    
    def _get_partial_result(self, session_id: str, step_name: str) -> dict:
        """部分的な結果を取得"""
        response = self.dynamodb.query(
            TableName=self.table,
            KeyConditionExpression='sessionId = :sid',
            ExpressionAttributeValues={':sid': {'S': session_id}}
        )
        
        # 最新の部分結果を返す
        items = response.get('Items', [])
        if items:
            return {'partial': True, 'lastSavedState': items[-1]}
        return {}
    
    def _send_alert(self, session_id: str, failed_step: str, error: Exception, completed_steps: list):
        """失敗アラートを送信"""
        message = f"""
Bedrock Agent実行エラー
- Session ID: {session_id}
- Failed Step: {failed_step}
- Error: {str(error)}
- Completed Steps: {len(completed_steps)}
- Completed: {[s['step'] for s in completed_steps]}
"""
        
        self.sns.publish(
            TopicArn=self.topic,
            Subject=f"Bedrock Agent Error: {failed_step}",
            Message=message
        )
    
    def _execute_single_step(self, session_id: str, step: dict) -> dict:
        # ステップ実行のダミー実装
        return {'success': True}

2026年時点での最新知見:プロンプトとメモリ最適化

2026年の現在、Bedrock Agentsの環境も随分進化してます。特に重要な変化が2つあります。

1. プロンプトキャッシング

Bedrock Agentsにプロンプトキャッシング機能が統合されたので、マルチステップでも同じシステムプロンプトを何度も処理しなくて済むようになりました。これだけで20~30%のトークン削減になるんですよね。正直、導入前と後で処理速度が全然違います。

2. 分散セッション状態

複数のエージェントインスタンスで同じセッションを処理する場合、DynamoDBのGSI(グローバルセカンダリインデックス)を活用して状態を効率的に共有できるようになってます。スケーラビリティの面で大幅に改善されました。

まとめ

Bedrock Agentsのマルチステップ設計で本番環境を安定運用するには、以下の3点が本当に重要です。

1. メモリ管理が全て — これは本当に強調したいんですが、コンテキストサイズを常に監視して、不要なデータはアーカイブするしかありません。50KB制限を徹底することで、ステップ数が増えても応答時間が線形に伸びる問題を抑止できます。逆にこれをしないと100%本番で落ちます。

2. ツール呼び出しに上限を設ける — 無限ループはほぼ必ず起きます。ステップごとに最大呼び出し回数とタイムアウトを設定し、自動で強制終了する仕組みが必須です。テスト環境では出なかった問題が本番だと突然出るんですよ。

3. 状態遷移とエラー復旧戦略 — ステップ間で大きなデータを引き継ぐときはS3を活用して、失敗時は部分結果で進める柔軟性を持たせることが、本番安定性を大きく向上させます。ここまで考え抜くことで、初めて運用に耐える構成になるんだと思います。

最初は「複数エージェント構成は高度で難しい」と思ってたんですが、むしろ設計思想をシンプルに保つことが正解だったんですよ。各ステップは単一責任で、状態は最小限に、エラーは想定しておく——その繰り返しです。毎回同じことの繰り返しなんですけど、それが安定運用を生み出すんですね。

U

Untanbaby

ソフトウェアエンジニア|AWS / クラウドアーキテクチャ / DevOps

10年以上のIT実務経験をもとに、現場で使える技術情報を発信しています。 記事の誤りや改善点があればお問い合わせからお気軽にご連絡ください。

関連記事