Database

Amazon DynamoDB Streams

DynamoDBの項目レベル変更を時系列で流す仕組み。AIP-C01では「業務データが更新されたら、その変化をRAGの検索対象に近似リアルタイムで反映する」配管として出ます。

  • B|実装まわり
  • D1 FM統合・データ・コンプラ
  • D2 実装と統合

ひとことで言うと

DynamoDBテーブルへの変更(作成・更新・削除)を、24時間だけ保持する時系列ログとして流す仕組みです。単体では何もしません。誰かが読んで処理して初めて意味を持つ配管です。

要するに、レジの取引ログです。棚(テーブル)の在庫が動くたびに、何がどう変わったかを控えのテープに刻む。テープ自体は何もしませんが、そのテープをLambdaやOpenSearch Ingestionに読ませれば、「在庫が変わったら検索インデックスも即座に更新する」が組めます。

なぜ生まれたか

RAGで社内データを使う場合、元のDynamoDBテーブルが更新されたのに、検索側(ベクトルストア)が古いままという乖離が起きます。バッチで毎晩同期する設計では、更新から反映までのタイムラグが読者体験を損ないます。

起きる不便Streamsでの解決
元データの更新が検索結果に反映されるまで遅い変更をほぼリアルタイムでLambda等に配信
更新前後の値を比較したい(差分検知)StreamViewTypeでNEW_AND_OLD_IMAGESを選べる
TTLによる自動削除と、ユーザーによる削除を区別したいTTL削除は「サービスによる削除」として区別されてストリームに現れる

何に使うか

用途使い方関連ドメイン
RAG対象データの近似リアルタイム同期OpenSearch ServiceへのZero-ETL統合の差分連携部分として利用D1
変更をトリガに埋め込み再生成Lambdaトリガーで変更項目だけをBedrockの埋め込みモデルに送るD1・D2
会話履歴の削除監査TTLによる自動削除をストリームで検知し、監査ログに記録D3

試験でどう問われるか

問われ方正解に寄る条件引っかけの選択肢
商品カタログの更新をRAGの検索結果に速く反映したいDynamoDB StreamsベースのZero-ETL統合でOpenSearchへ差分同期毎晩バッチでKnowledge Baseを再同期する案(遅い)
変更前後の値を比較して処理を分岐したいStreamViewType: NEW_AND_OLD_IMAGESKEYS_ONLYを選ぶ(値の中身が取れない)
24時間より前の変更履歴を後から調べたいStreamsの保持期間は24時間限定。長期保持が要るならKinesis Data Streams等への転送を別途設計Streamsだけで長期保存できると誤認する選択肢

保持期間は24時間。「無効化しても24時間は読める」「24時間を過ぎたら自動的にトリムされる」の両方が問われます。長期保持・再生が要件に出たら、Streams単体では答えになりません。

実務でどう使うか

  • Lambdaのイベントソースマッピングはシャード単位で処理される。 1シャードは既定で1インスタンスが順序を守って読みます。並列度を上げたいときはParallelizationFactor(最大10)で、同一アイテムの順序を保ったまま複数インスタンスに分散できます
  • 同じシャードを2プロセス以上で読むとスロットリングされます。 独自にコンシューマを複数走らせる設計は避ける
  • 値が変わらないPutItem/UpdateItemはストリームに記録されません。 「更新のたびに埋め込みを作り直す」設計で、無駄な再計算を避けられる一方、「更新したのに反映されない」場合はここを疑う

ハマりどころ: Zero-ETL統合の初回同期はPoint-in-Time RecoveryによるS3エクスポートが前提で、Streamsが担うのはそれ以降の差分だけです。PITRを有効化していないと初回のフルスナップショットが作れません。

取り違えやすいもの

迷う相手切り分けの一言
DynamoDB本体本体=データを置く場所。Streams=その変化を流す配管。単独では検索にも通知にも使えない
Kinesis Data StreamsDynamoDB Streamsは24時間限定・DynamoDB専用。長期保持や複数コンシューマでの再生が要るならKinesis Data Streams用のアダプタ経由で転送する設計になる
EventBridgeStreamsは項目レベルの変更配信。EventBridgeはアプリ間のイベントルーティング。「DynamoDBの変更を他サービスに広く配りたい」ならStreams→Lambda→EventBridgeの組み合わせになる

想起チェック

Q1. 商品カタログ(DynamoDB)の更新を、RAGの検索結果にできるだけ早く反映したい。何を使いますか?

DynamoDB StreamsをベースにしたOpenSearch ServiceへのZero-ETL統合です。初回はPITRによるS3エクスポート、以降はStreamsによる差分同期で近似リアルタイムに反映します。

Q2. 「1週間前の変更履歴を調べたい」と言われたら、Streamsだけで答えられますか?

答えられません。保持期間は24時間限定です。長期の変更履歴が必要なら、Kinesis Data StreamsやS3への転送を別途設計する必要があります。

出典(AWS公式)