← AWSサービスの内部原理 コース
52. パーティションキーとシャード — MD5ハッシュがストリームを水平分割する
Kinesisのストリームがどうやって1本の流れを複数のシャードに分け、なぜ「同じパーティションキーだけ」が順序を持つのかを、レコードが入る瞬間から図で追っていきます。
① レコードはパーティションキーを持って入ってくる
プロデューサーputRecord / putRecords
内訳
データレコード1レコード = 3要素
パーティションキーUnicode文字列・最大256文字
データ本体(blob)最大1MB・中身は非解釈
シーケンス番号書き込んだ後にKinesisが付与
入口ではまだ無い- 公式の言い回し:「Data records are composed of a sequence number, a partition key, and a data blob」。
- 「Partition keys are Unicode strings, with a maximum length limit of 256 characters for each key.」
- 「A data blob can be up to 1 MB.」データ本体は「does not inspect, interpret, or change」= Kinesisは中身を見ない。
- シーケンス番号は「Kinesis Data Streams assigns the sequence number after you write to the stream」= 入口ではまだ決まっていない(④で効いてくる)。
レコードが持ち込むのはパーティションキーとデータ本体だけ。どのシャードに入るかも、順序を決めるシーケンス番号も、この時点ではまだ決まっていない。
② MD5でキーを128bit整数に変換し、シャードの範囲に割り当てる
パーティションキー文字列
どの範囲に入るか?
128bit整数ハッシュキー
同じキー → 必ず同じハッシュキー各シャードが連続した範囲を1つずつ担当
ハッシュキー空間0〜2^128-1 を分割
シャードA範囲 0..k1
シャードB範囲 k1..k2
シャードC範囲 k2..max
- 公式:「An MD5 hash function is used to map partition keys to 128-bit integer values and to map associated data records to shards using the hash key ranges of the shards.」
- 「Data records that have the same partition key also have the same hash key value.」= 同じキーは必ず同じシャードへ。
- 各シャードの範囲は「a set of ordered contiguous non-negative integers」で、StartingHashKey 〜 EndingHashKey で表される。
- 「0〜2^128-1」という上限値そのものは、MD5が128bitであることからの導出(公式ドキュメントは範囲を「連続した非負整数の集合」と表現している)。
これがハッシュパーティショニング。キーをハッシュして数値にし、その数値を「どのシャードの範囲に入るか」で振り分ける。DynamoDB編で見たパーティションキーと同じ原理の別実装で、キー→ハッシュ→担当ノード(シャード)という流れは共通している。
③ シャードは容量の固定単位 — ストリーム容量はその和
1シャード固定の容量ユニット
足し算するだけ
書き込み最大1MB/秒・1,000レコード/秒
MBにはパーティションキーも含む読み取り最大2MB/秒・5トランザクション/秒
容量が足りない/余る?
ストリーム容量= 各シャード容量の総和
例: 3シャード → 書き込み3MB/秒・3,000rec/秒シャード数を増減= ⑥リシャーディング
- 公式:「Each shard can support up to 5 transactions per second for reads, up to a maximum total data read rate of 2 MB per second and up to 1,000 records per second for writes, up to a maximum total data write rate of 1 MB per second (including partition keys).」
- 「The total capacity of the stream is the sum of the capacities of its shards.」
- provisionedモードでは自分でシャード数を指定する(「you must specify the number of shards for the data stream」)。
シャードは「速さ」の単位。1つあたりの上限が決まっていて、ストリームの容量はその単純な足し算。だから容量を変えたければシャードの数そのものを変える(⑥)。
④ 順序が保証されるのはシャードの中だけ
シャードの中シーケンス番号が振られる
同一シャード内・同一キーで一意/通常(generally)増加シャードをまたぐと?
#1 → #2 → #3順序が保たれる
だから…
シャードAの #x前後関係なし
シャードBの #y番号は横断インデックスにできない
同じキーにまとめる順序が欲しいレコードは
- 公式:「Each data record has a sequence number that is unique per partition-key within its shard.」
- 「Sequence numbers for the same partition key generally increase over time.」= 厳密な連番ではなく「概ね増加」。書き込み間隔が長いほど番号の開きも大きくなる。
- Note:「Sequence numbers cannot be used as indexes to sets of data within the same stream.」= シャードをまたいだ順序付けには使えない。
順序保証のスコープはシャード内(実質は同じパーティションキー内)に閉じている。「順序のスコープをキーで区切る」という考え方は、SQS FIFOのメッセージグループIDと同型 — 同じグループ/同じキーの中だけが順序を持ち、異なるグループ同士は独立して並列に流れる。
⑤ 偏るとホットシャードになる — キー設計が効く理由
悪いキー設計少数のキーに集中(例: 大半が "JP")
Aだけ上限(1MB/s・1,000rec/s)に張り付く
シャードA●●●●●●●●●● 満杯
シャードB・ 空き
シャードC・ 空き
対策: カーディナリティの高いキー(例: ユーザーID)
ホットシャード他が空いていてもスループット頭打ち
シャードA●●● 均等
シャードB●●● 均等
シャードC●●● 均等
- ③の上限は「シャードごと」に効く。ストリーム全体に余裕があっても、1シャードに負荷が偏れば、そのシャードの1MB/s・1,000rec/sで律速する。
- MD5は決定的なので同じキーは必ず同じシャードに落ちる(②)。偏りはハッシュのランダム性ではなくキーの選び方で決まる。
容量が「シャードごとの固定単位」だからこそ、キーが偏ると特定シャードだけが満杯になる=ホットシャード。DynamoDB編のホットパーティションと同じ構図で、対策も同じ「カーディナリティの高いキーで負荷を散らす」。
⑥ 容量変更=リシャーディング — 常にペアで、親を読み切ってから子へ
リシャーディング常にペアワイズ
1回で3つ以上には割れない/結合できないスプリットでは親のハッシュ範囲を子で分け合う
スプリット1つの親 → 2つの子(容量↑)
マージ2つの親 → 1つの子(容量↓)
指定した値 v を境に
親の範囲[ start ... end ]
親はどうなる?
子1[ start .. v )
子2[ v .. end ]
子から先に読むと同じキーの順序が崩れる
親の状態遷移OPEN → CLOSED → EXPIRED
CLOSED: 新規は子へ・既存レコードはまだ親に残るこの一連の調整を自動化したのが
親を読み切ってから子へgetNextShardIterator が null = 読み切った合図
オンデマンドモード必要なスループットに合わせシャード管理を自動で行う
- 公式:「Resharding is always pairwise in the sense that you cannot split into more than two shards in a single operation, and you cannot merge more than two shards in a single operation.」親=parent、子=child。
- スプリット:「That hash key value and all higher hash key values are distributed to one of the child shards. All the lower hash key values are distributed to the other child shard.」= 親の範囲を子で分け合う。分割点は範囲内の任意の値でよい(半分にするのは一例)。
- 「If you read data from the child shards before having read all data from the parent shards, you could read data for a particular hash key out of the order given by the data records' sequence numbers. ... you should, after a reshard, always continue to read data from the parent shards until it is exhausted.」
- 「When getRecordsResult.getNextShardIterator returns null, it indicates that you have read all the data in the parent shard.」
- 親の状態: OPEN(読み書き可)→ CLOSED(書き込みは子へ、読み取りは保持期間内の限られた時間 "for a limited time" だけ可)→ EXPIRED(保持期間切れで読めなくなる)。
- オンデマンド:「Kinesis Data Streams automatically manages the shards in order to provide the necessary throughput.」なお SplitShard/MergeShards APIはprovisionedモード専用で、オンデマンドのストリームに対して呼ぶとValidationExceptionになる=オンデマンドではこの調整をユーザーが行わない。
容量変更はシャードの分割/マージで行い、常に親→子のペアでハッシュ範囲を受け渡す。順序を壊さない鍵は「親を読み切ってから子へ」という順番で、これは④のシーケンス番号がシャード内に閉じているという性質の当然の帰結。オンデマンドモードはこの一連の調整を自動化したもの。