ひとことで言うと
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_IMAGES | KEYS_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 Streams | DynamoDB Streamsは24時間限定・DynamoDB専用。長期保持や複数コンシューマでの再生が要るならKinesis Data Streams用のアダプタ経由で転送する設計になる |
| EventBridge | Streamsは項目レベルの変更配信。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への転送を別途設計する必要があります。