ハイパーグロースのための最新データベースシャーディング戦略

Table of Contents
データベースのスケーリングにおいて、シャーディングは最終手段です。シャーディングは、望んで行うものではなく、プライマリライターが飽和状態になり、垂直スケーリングが物理的な限界(またはAWSの請求上限)に達し、リードレプリカでは書き込みスループットのボトルネックを解決できない場合に実施するものです。
この限界を超えると、単一ノードのリレーショナルデータベースが提供する保証(ACIDトランザクション、テーブル間の外部キー、グローバルな一意制約)を犠牲にして、水平スケーラビリティを獲得することになります。このガイドでは、現代のシステムがどのようにリレーショナルデータベースをシャーディングしているのか、シャーディングキーの選び方、そしてそれに伴う分散システムの障害モードへの対処法について解説します。
シャーディングすべきタイミング(そして最初に試すべきこと)
データベースをシャードに分割する前に、以下の暫定的なステップをすべて試したことを確認してください。
Scaling Stages:
1. Index Optimization & Slow Query Tuning (EXPLAIN ANALYZE)
│ (hit CPU/IOPS limits)
▼
2. Read/Write Splitting (Primary for writes, Replicas for reads)
│ (write throughput saturates primary node)
▼
3. Table Partitioning (PostgreSQL declarative partitioning / MySQL partition)
│ (single machine disk/memory/lock contention saturated)
▼
4. Functional Partitioning (Extract Auth, Billing, Analytics to isolated DBs)
│ (a single domain table like `orders` or `events` exceeds 1 node capacity)
▼
5. Horizontal Sharding (Split identical table schema across N database nodes)
単一のテーブルが5億行以上を超え、書き込み操作が継続的なロック競合を引き起こし、WAL/redoレプリケーションの遅延がレプリカが取り込める速度よりも速く増加する場合、水平シャーディングが必要です。
1. シャーディングキー戦略
シャーディングキーは、特定の行がどのデータベースシャードに格納されるかを決定します。誤ったキーを選択すると、ホットスポット、クロスシャードの散発的なクエリ、そして悪夢のようなマイグレーションにつながります。
戦略A:範囲ベースのシャーディング
行は連続する値の範囲(例:created_at日付範囲やID範囲1–1,000,000 -> シャード1、1,000,001–2,000,000 -> シャード2)に基づいてルーティングされます。
- 利点: 単一シャード内での日付/ID範囲のクエリが容易。シャード全体を削除することで古いデータをアーカイブするのが簡単。
- 欠点(主要): 深刻な書き込みホットスポット。すべての新しい書き込みは最も高い範囲のシャード(例:今日の日付)に集中し、古いシャードはアイドル状態になる一方で、最新のノードは過負荷になります。
戦略B:ハッシュベースのシャーディング(コンシステントハッシュ)
シャーディングキーを暗号学的または高速な非暗号学的ハッシュ関数(MurmurHash3やxxHashなど)に通し、シャード数で割った余りを取るか、仮想ノードを持つコンシステントハッシュリングにマッピングします。
\text{Shard ID} = \text{hash}(\text{sharding\_key}) \pmod N
import mmh3
class ShardRouter:
def __init__(self, shard_count: int):
self.shard_count = shard_count
def get_shard_id(self, entity_id: str) -> int:
# MurmurHash3 ensures uniform distribution across shards
hash_val = mmh3.hash(str(entity_id))
return abs(hash_val) % self.shard_count
router = ShardRouter(shard_count=8)
print(router.get_shard_id("user_98314")) # -> Shard 3
print(router.get_shard_id("user_98315")) # -> Shard 7
- 利点: すべてのノードに読み書きが均等に分散される。書き込みホットスポットがない。
- 欠点: キーをまたぐ範囲クエリ(例:
WHERE created_at BETWEEN x AND y)は、すべてのシャードにブロードキャストする必要がある(散発的なクエリ)。
戦略C:ディレクトリ/ルックアップベースのシャーディング
中央の高度にキャッシュされたルックアップサービス(例:Redis + メタデータデータベース)がentity_idをshard_idにマッピングします。
- 利点: 最高の柔軟性。ルックアップテーブルを更新するだけで、個々の大量の顧客(「クジラ」テナント)を専用のハードウェアシャードに移動できます。
- 欠点: すべてのデータベース操作にルックアップサービスへの追加のネットワークラウンドトリップが発生する。ルックアップサービスが十分にキャッシュされていない場合、単一障害点(SPOF)となる。
2. 実装例:Pythonでのアプリケーションレベルルーティング
アプリケーションレベルのシャーディングでは、サービス層がシャードごとのコネクションプールを管理し、コンテキストに基づいてクエリをルーティングします。
import asyncpg
from typing import Dict, Any
class ShardedDatabaseManager:
def __init__(self, shard_configs: Dict[int, str]):
self.configs = shard_configs
self.pools: Dict[int, asyncpg.Pool] = {}
async def initialize(self):
for shard_id, dsn in self.configs.items():
self.pools[shard_id] = await asyncpg.create_pool(
dsn, min_size=5, max_size=20, timeout=10.0
)
def get_shard_for_tenant(self, tenant_id: str) -> int:
import hashlib
h = int(hashlib.md5(tenant_id.encode()).hexdigest(), 16)
return h % len(self.pools)
async def execute_query(
self, tenant_id: str, query: str, *args
) -> list[Dict[str, Any]]:
shard_id = self.get_shard_for_tenant(tenant_id)
pool = self.pools[shard_id]
async with pool.acquire() as conn:
records = await conn.fetch(query, *args)
return [dict(r) for r in records]
async def scatter_gather(self, query: str, *args) -> list[Dict[str, Any]]:
"""Broadcast read query to all shards concurrently and aggregate results."""
import asyncio
async def query_shard(shard_id: int, pool: asyncpg.Pool):
async with pool.acquire() as conn:
return await conn.fetch(query, *args)
tasks = [query_shard(sid, pool) for sid, pool in self.pools.items()]
results = await asyncio.gather(*tasks, return_exceptions=False)
flat = []
for r in results:
flat.extend([dict(row) for row in r])
return flat
3. シャーディングされたリレーショナルデータベースの4つの大きな課題
1. グローバルな一意識別子の問題
各シャードが衝突するプライマリキーを生成するため、AUTO_INCREMENTやPostgreSQLのBIGSERIALシーケンスに頼ることはできません。
解決策:
- UUIDv7: 時間順の128ビットUUID(データベースインデックスの局所性が良好で、中央での調整が不要)。
- Snowflake ID:
[Timestamp (41 bits) | Worker ID (10 bits) | Sequence (12 bits)]で構成される64ビット整数。データベースロックなしで、毎秒数百万個のソート可能な一意IDを生成します。
Snowflake 64-Bit Structure:
[ 1 bit sign ] [ 41 bits timestamp (ms) ] [ 10 bits machine/shard ID ] [ 12 bits sequence ]
2. クロスシャード結合
もしordersがtenant_idでシャードされ、productsがグローバルである場合、異なる物理サーバー間でSQLで直接結合することは不可能です。
- パターンA: エンティティグルーピング(コロケーション)。 すべての関連データを同じルートキーでシャードします。もし
users、orders、およびorder_itemsがすべてtenant_idをシャーディングキーとして含んでいる場合、単一のtenant_idに限定されたクエリは、その単一シャード上で100%ローカルな結合として実行されます。 - パターンB: 参照テーブルのレプリケーション。 小規模でめったに変更されないルックアップテーブル(例:
currencies、countries、plan_tiers)は、すべてのシャードにコピーされます。更新は、変更データキャプチャ(Debezium/Kafka)を介してすべてのノードにブロードキャストされます。 - パターンC: アプリケーションレベルのインメモリ結合。 シャードAから
ordersをクエリし、一意のproduct_idsを取得し、カタログデータベースに対してバッチクエリSELECT * FROM products WHERE id = ANY(...)を実行し、バックエンドサービスでDTOを結合します。
3. 分散トランザクション:2PC vs Saga
ビジネス操作が2つの異なるシャードにまたがる場合(例:シャード1の口座からシャード2の口座への送金)、標準のBEGIN ... COMMITではアトミック性を保証できません。
Two-Phase Commit (2PC)
トランザクションコーディネーターは、参加するすべてのシャードに「コミットできますか?」と尋ねます(準備フェーズ)。すべてがYESと返答した場合、コーディネーターはコミットをログに記録し、「コミット!」を発行します(コミットフェーズ)。
- 問題点: 高いレイテンシ、ネットワークホップをまたぐロック保持、コーディネーター障害によるブロッキングモード。高スループットのインターネットシステムではほとんど使用されません。
トランザクションアウトボックスを用いたSagaパターン
同期的な分散ロックを非同期的な補償トランザクションに置き換えます。
1. [Shard 1] Deduct $100 from Account A -> Write "MoneyDebited" event to local `outbox` table (Atomic local commit)
2. CDC worker reads `outbox` -> Publishes to Kafka
3. Worker reads Kafka -> [Shard 2] Credit $100 to Account B
4. If Shard 2 fails permanently -> Emit "CreditFailed" event -> Compensating transaction on Shard 1 to refund $100 to Account A
4. 最新のデータベースミドルウェア vs ネイティブ分散SQL
アプリケーションコードでカスタムルーティング層を維持する代わりに、最新のアーキテクチャではプロキシまたは分散ストレージ層に依存することがよくあります。
| アプローチ | 例 | 仕組み | 最適な用途 |
|---|---|---|---|
| 透過的プロキシ | Vitess (MySQL), Citus (Postgres) | アプリケーションと標準のMySQL/Postgresの間に位置する。アプリは標準SQLを話し、プロキシがASTを解析してシャードをルーティングする。 | アプリケーションSQLを書き換えることなく、既存のモノリシックなリレーショナルデータベースをスケーリングする場合。 |
| 分散SQL(ネイティブ) | CockroachDB, TiDB, Google Spanner | ストレージエンジンがネイティブに分散されている(Raftコンセンサス + LSMツリー/Pebble)。シャード(範囲)は自動的に分割され、再バランスされる。 | 標準ACIDトランザクションを備えた水平書き込みが必要な、グリーンフィールドのクラウドネイティブシステム。 |
| アプリレベルシャーディング | カスタムルーター / Hibernate Shards / SQLAlchemy Routing | アプリケーションがエンティティIDに基づいて接続文字列を明示的に選択する。 | プロキシのレイテンシオーバーヘッド(1~2ms)が許容できない、極めて高いスループットの場合。 |
5. ライブリバランス:オンラインマイグレーションのレシピ
いずれ、8つのシャードでは容量が不足し、プラットフォームをオフラインにすることなく16のシャードに拡張する必要が生じます。
- 新しいシャードノードを起動します(シャード9~16)。
- デュアルライト: アプリケーションを更新し、古いシャードの場所と新しいシャードの場所の両方に書き込むようにします。
- バックフィル: バックグラウンドワーカーを実行し、古いシャードから新しいシャードに履歴データをコピーします(新しいコンシステントハッシュルーティングでフィルタリング)。
- 整合性検証: 古いシャードと新しいシャード間でチェックサムを比較します。
- リードの切り替え: 読み取りトラフィックを新しい16シャードマッピングに振り向けます。
- デュアルライトの停止: シャード1~8から移行済みの古いレコードを削除します。
シャーディングアーキテクチャの要約チェックリスト
- シャーディングキーを慎重に選択する: テナント/ユーザーごとにデータをコロケーションし、95%以上のクエリを単一シャードに保つ。
- ノード間のプライマリキーIDの衝突を避けるために、64ビットのSnowflakeまたはUUIDv7を使用する。
- Sagaパターンとトランザクションアウトボックスを使用して、クロスシャードトランザクションを排除する。
- ローカル結合を可能にするために、小さな参照テーブルをすべてのシャードにレプリケートする。
- 数百行のカスタムアプリケーションレベルルーティングロジックを記述する前に、Vitess / Citusをベンチマークする。
こちらもおすすめです
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles
PostgreSQLのVACUUMとインデックス肥大化:検知、軽減、そして自動チューニング
PostgreSQLのテーブルとインデックスの肥大化を診断・解消します。自動バキュームのチューニング方法、pg_repackによるゼロダウンタイムでの再構築、MVCCの可視性マップまでを解説します。
Read more
pgvectorとハイブリッド検索でプロダクションRAGを構築する
pgvectorとハイブリッド検索を組み合わせた堅牢なRAGアーキテクチャの構築方法を学び、全文ベクトル検索とBM25を組み合わせて検索精度を向上させ、HNSWインデックスを調整し、プロダクションPythonクライアントを出荷します。
Read more
Redisメモリ最適化:内部構造、データ構造エンコーディング、メモリプロファイリング
RedisのRAM使用量を最大70%削減し、ziplists、listpacks、quicklists、string SDSのオーバーヘッド、自動メモリ断片化軽減策について深く掘り下げます。
Read more