Kafkaの階層型ストレージアーキテクチャ: AWS S3とGoogle Cloud StorageでEBSストレージコストを削減

目次(19 項目)
Apache Kafkaの従来のアーキテクチャでは、コンピューティングとストレージが密接に結合されています。ブローカーはローカルのログセグメントを管理し、通常はAWS EBSやGCP Persistent Diskのような高性能ブロックストレージによってバックアップされています。この設計は、最近のデータに対しては比類のないスループットと低遅延アクセスを提供しますが、特にデータ量がペタバイト規模に拡大すると、長期的なデータ保持には経済的に持続不可能になります。履歴データやアクセス頻度の低いデータに対して、ブロックデバイスのプロビジョニングされたIOPSとストレージ容量のコストは、インフラ予算をすぐに圧迫する可能性があります。
KIP-405として導入されたKafka Tiered Storageは、Kafkaのストレージ層を根本的に再構築します。これにより、古くアクセス頻度の低いログセグメントを、ローカルのブローカーストレージから、AWS S3やGoogle Cloud Storage(GCS)のような費用対効果が高く、耐久性の高いオブジェクトストレージソリューションにシームレスにオフロードできます。このコンピューティングとストレージの分離により、Kafkaのコア保証を損なうことなく、大幅なコスト削減、スケーラビリティの向上、運用管理の簡素化が可能になります。
Kafka Tiered Storage (KIP-405) の理解
KIP-405は、2層のストレージモデルを導入しています。Kafkaブローカー上のローカル層と、オブジェクトストレージ上のリモート層です。主な目的は、ホットデータに対してローカルストレージのパフォーマンス上の利点を維持しつつ、履歴データに対してクラウドオブジェクトストレージの費用対効果と実質的に無限のスケーラビリティを活用することです。
コアコンセプト
- コンピューティングとストレージの分離: ブローカーは、すべてのトピックデータを保存する責任を単独で負うことはなくなります。主にローカルディスクから最新のデータを提供し、古いデータをリモート層にオフロードします。これにより、コンピューティング(ブローカー)とストレージ(オブジェクトストレージ)を独立してスケーリングできます。
- ローカルストレージ層: これは、ブローカーのローカルディスク上の従来のKafkaログディレクトリです。最新のログセグメントを保持し、アクティブなプロデューサーとコンシューマーに低遅延の読み書きアクセスを提供します。この層の保持期間は設定可能です。
- リモートストレージ層: これはクラウドオブジェクトストレージ(S3、GCS)です。設定された保持ポリシーに基づいてログセグメントが「コールド」と判断されると、非同期的にこの層にアップロードされます。この層は、すべての履歴データに対する不変で耐久性の高いアーカイブとして機能します。
- リモートストレージマネージャー: Kafkaブローカー内のコンポーネントで、ローカルストレージとリモートストレージ間のログセグメントのライフサイクルを管理します。セグメントのオブジェクトストレージへのアップロードと、コンシューマーがローカル保持期間を超えたデータを要求した際のセグメントの取得を処理します。
- メタデータ管理: 両方の層からのシームレスな消費を可能にするために、Kafkaはどのセグメントがローカルに存在し、どのセグメントがリモートストレージにオフロードされたかに関するメタデータを維持します。このメタデータは、コンシューマーが物理的な場所に関係なくデータを透過的にフェッチするために不可欠です。
階層型ストレージの利点
- 大幅なコスト削減: オブジェクトストレージは、長期保持の場合、ブロックストレージ(EBS/Persistent Disk)よりも桁違いに安価です。これが導入の主な動機です。
- スケーラビリティの向上: ストレージをブローカーから分離することで、ブローカーを追加したりローカルディスクサイズを増やしたりすることなく、ストレージ容量をほぼ無限にスケーリングできます。
- 運用の簡素化: ブローカーは、より小さく高速なローカルディスクでプロビジョニングできるため、障害発生時の回復時間が短縮され、ディスク管理が簡素化されます。
- データ保持期間の延長: 法外なコストをかけることなく、数ヶ月または数年間データを保持できるため、新しい分析ユースケースやコンプライアンス要件に対応できます。
- ブローカーの回復性の向上: ローカルディスクのフットプリントが小さいほど、ブローカー障害発生時のデータレプリケーションと回復が高速になります。
アーキテクチャの詳細
階層型ストレージアーキテクチャは、Kafkaがログセグメントを管理する方法を根本的に変更します。データフローとコンポーネントの相互作用を理解することは、効果的なデプロイと最適化のために不可欠です。
データフローとコンポーネントの相互作用
- プロデューサーの書き込み: プロデューサーはメッセージをトピックパーティションに追加します。これらのメッセージは、ブローカーのローカルディスク上のアクティブなログセグメントに書き込まれます。
- セグメントのロールオーバー: アクティブなログセグメントがサイズ(
log.segment.bytes)または時間(log.segment.ms)の制限に達すると、ロールオーバーされ、不変のクローズドセグメントになります。 - ローカル保持ポリシー: ブローカーは、ローカルに保持するセグメントを決定するために
local.retention.bytesまたはlocal.retention.msを適用します。このポリシーよりも古いセグメントは、オフロード対象としてマークされます。 - リモートセグメントのアップロード:
RemoteLogManagerは、これらのマークされたセグメントを設定されたオブジェクトストレージバケット(S3またはGCS)に非同期的にアップロードします。このプロセスは、ブローカーの操作をブロックしません。 - ローカルセグメントの削除: セグメントがリモートストレージに正常にアップロードされ、そのメタデータが更新されると、ローカルコピーを削除してローカルディスクスペースを解放できます。
- コンシューマーの読み取り:
- ホットデータ: コンシューマーがローカル保持期間内のデータを要求した場合、ブローカーはローカルディスクから直接データを提供し、低遅延を維持します。
- コールドデータ: コンシューマーがリモートストレージにオフロードされたデータを要求した場合、ブローカーはS3/GCSから必要なセグメントを透過的にフェッチし、一時的にキャッシュしてコンシューマーに提供します。これにより、追加の遅延が発生します。
主要な設定パラメータ
remote.log.storage.enable: 階層型ストレージを有効/無効にするグローバルブローカー設定。trueである必要があります。remote.log.storage.manager.impl.class: リモートストレージマネージャーの実装クラスを指定します。- AWS S3の場合:
org.apache.kafka.server.log.remote.metadata.storage.S3RemoteLogMetadataManager - Google Cloud Storageの場合:
org.apache.kafka.server.log.remote.metadata.storage.GcsRemoteLogMetadataManager
- AWS S3の場合:
remote.log.storage.manager.impl.prefix: すべてのリモートストレージマネージャー関連設定のプレフィックス。例:remote.log.storage.manager.impl.aws.region。local.retention.bytes: パーティションごとにローカルディスクに保持する最大バイト数。この制限を超えると、古いセグメントがオフロードされます。local.retention.ms: パーティションごとにローカルディスクにセグメントを保持する最大時間(ミリ秒)。log.segment.bytes: ログセグメントファイルの最大サイズ。セグメントが小さいほど、ロールオーバーとアップロードの頻度が高くなり、オブジェクトストレージAPI呼び出しが増加する可能性があります。log.segment.ms: ログセグメントがロールオーバーされるまでの最大時間。
設定と実装
Kafka Tiered Storageを実装するには、ブローカーレベルとトピックレベルの両方で慎重な設定が必要です。AWS S3とGoogle Cloud Storageのセットアップについて詳しく説明します。
前提条件
- Kafkaバージョン: KIP-405にはApache Kafka 3.0.0以降が必要です。本番環境では、安定性と機能のために3.3.x以降が推奨されます。
- クラウド認証情報:
- AWS S3: ターゲットS3バケットに対する
s3:PutObject、s3:GetObject、s3:DeleteObject、s3:ListBucket、s3:GetBucketLocationの権限を持つIAMロールまたはユーザー。 - Google Cloud Storage: ターゲットGCSバケットに対する
Storage Object AdminまたはStorage Object CreatorおよびStorage Object Viewerロールを持つサービスアカウントキー。
- AWS S3: ターゲットS3バケットに対する
- クラウドバケット: 目的のリージョンに作成されたS3バケットまたはGCSバケット。
ブローカー設定 (server.properties)
以下のパラメータは、各Kafkaブローカーのserver.propertiesに追加されます。
共通の階層型ストレージ設定
# Enable remote storage
remote.log.storage.enable=true
# Set the local retention policy.
# Example: Retain 10GB per partition locally.
# This value should be carefully chosen based on hot data access patterns.
local.retention.bytes=10737418240 # 10 GB
# Alternatively, or in addition to bytes, set time-based local retention.
# Example: Retain 24 hours of data locally.
# local.retention.ms=86400000 # 24 hours
# Configure the segment size and time for optimal offloading.
# Smaller segments mean more frequent uploads but faster local deletion.
log.segment.bytes=1073741824 # 1 GB
log.segment.ms=604800000 # 7 days (default)
AWS S3固有の設定
# S3 Remote Log Storage Manager implementation class
remote.log.storage.manager.impl.class=org.apache.kafka.server.log.remote.metadata.storage.S3RemoteLogMetadataManager
# S3 bucket name for remote storage
remote.log.storage.manager.impl.aws.bucket.name=your-kafka-tiered-storage-s3-bucket
# AWS region of the S3 bucket
remote.log.storage.manager.impl.aws.region=us-east-1
# (Optional) AWS endpoint for S3, useful for S3-compatible storage or private endpoints
# remote.log.storage.manager.impl.aws.endpoint=https://s3.us-east-1.amazonaws.com
# (Optional) AWS credentials provider. Default is DefaultAWSCredentialsProviderChain.
# remote.log.storage.manager.impl.aws.credentials.provider.class=com.amazonaws.auth.DefaultAWSCredentialsProviderChain
# (Optional) Number of threads for S3 uploads/downloads
# remote.log.storage.manager.impl.aws.num.client.threads=10
Google Cloud Storage固有の設定
# GCS Remote Log Storage Manager implementation class
remote.log.storage.manager.impl.class=org.apache.kafka.server.log.remote.metadata.storage.GcsRemoteLogMetadataManager
# GCS bucket name for remote storage
remote.log.storage.manager.impl.gcp.bucket.name=your-kafka-tiered-storage-gcs-bucket
# GCP Project ID
remote.log.storage.manager.impl.gcp.project.id=your-gcp-project-id
# Path to the service account key file (JSON).
# Ensure this file is accessible by the Kafka broker process.
# Alternatively, if running on GCP, use default credentials via instance service account.
# remote.log.storage.manager.impl.gcp.credentials.file=/path/to/your/gcp-service-account-key.json
# (Optional) Number of threads for GCS uploads/downloads
# remote.log.storage.manager.impl.gcp.num.client.threads=10
トピック設定
階層型ストレージは、ブローカーレベルのデフォルトを上書きして、トピックレベルで有効化および設定できます。
# Create a new topic with remote storage enabled and specific local retention
kafka-topics.sh --create --topic my-tiered-topic \
--bootstrap-server localhost:9092 \
--partitions 3 --replication-factor 3 \
--config remote.storage.enable=true \
--config local.retention.bytes=5368709120 # 5 GB local retention for this topic
# Modify an existing topic to enable remote storage
kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name my-existing-topic \
--alter --add-config remote.storage.enable=true
# Modify local retention for an existing tiered topic
kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name my-existing-topic \
--alter --add-config local.retention.ms=172800000 # 48 hours local retention
Strimzi を使用した Kubernetes デプロイメント
StrimziはKubernetes上でのKafkaの実行を簡素化します。階層型ストレージを有効にするには、Kafkaカスタムリソースを変更します。
Strimzi Kafka CR の例
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
name: my-kafka-cluster
spec:
kafka:
version: 3.6.0 # Ensure Kafka version supports KIP-405
replicas: 3
listeners:
- name: plain
port: 9092
type: internal
tls: false
- name: tls
port: 9093
type: internal
tls: true
storage:
type: jbod
volumes:
- id: 0
type: persistent-claim
size: 100Gi # Smaller local storage needed due to tiered storage
deleteClaim: false
config:
# Common Tiered Storage Settings
remote.log.storage.enable: "true"
local.retention.bytes: "10737418240" # 10 GB
log.segment.bytes: "1073741824" # 1 GB
# AWS S3 Specific Configuration
remote.log.storage.manager.impl.class: "org.apache.kafka.server.log.remote.metadata.storage.S3RemoteLogMetadataManager"
remote.log.storage.manager.impl.aws.bucket.name: "your-kafka-tiered-storage-s3-bucket"
remote.log.storage.manager.impl.aws.region: "us-east-1"
# For AWS, ensure the EC2 instance profile or Kubernetes service account has the necessary IAM role.
# Strimzi can inject environment variables for AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY if needed,
# but IAM roles are preferred.
# OR Google Cloud Storage Specific Configuration
# remote.log.storage.manager.impl.class: "org.apache.kafka.server.log.remote.metadata.storage.GcsRemoteLogMetadataManager"
# remote.log.storage.manager.impl.gcp.bucket.name: "your-kafka-tiered-storage-gcs-bucket"
# remote.log.storage.manager.impl.gcp.project.id: "your-gcp-project-id"
# For GCP, use Workload Identity to bind a Kubernetes Service Account to a GCP Service Account.
# The GCP Service Account should have Storage Object Admin permissions.
# Strimzi will automatically pick up credentials if the pod's service account is configured correctly.
# Example of how to configure a KafkaTopic with remote storage enabled
entityOperator:
topicOperator: {}
userOperator: {}
---
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
name: my-tiered-topic
labels:
strimzi.io/cluster: my-kafka-cluster
spec:
partitions: 3
replicas: 3
config:
remote.storage.enable: "true"
local.retention.bytes: "5368709120" # 5 GB local retention for this topic
Strimzi の IAM/GCP 権限
- AWS: Kubernetesノード(またはIRSA - IAM Roles for Service Accountsを使用している場合は特定のKafkaブローカーポッド)には、ターゲットS3バケットに対する
s3:PutObject、s3:GetObject、s3:DeleteObject、s3:ListBucket、s3:GetBucketLocation権限を付与するIAMロールがアタッチされている必要があります。 - GCP: Workload Identityを利用して、KubernetesサービスアカウントをGCPサービスアカウントにバインドします。GCPサービスアカウントには、GCSバケットに対する
Storage Object Adminまたは同等の権限が必要です。
コスト最適化分析
Kafka Tiered Storageの主な動機はコスト削減です。現実的なシナリオで潜在的な節約を定量化してみましょう。
シナリオ:
- 取り込みデータ: 10 TB/日
- レプリケーションファクター: 3(高可用性の標準)
- 総生データ: 30 TB/日
- 保持ポリシー: 「ホット」データは30日間、「コールド」データは1年間。
- 平均ブローカー数: 10ブローカー。
- リージョン: US East (N. Virginia) - AWSは
us-east-1、GCPはus-east4。
従来のKafka (EBSバックアップ)
1年間の保持の場合、必要なストレージの合計は30 TB/日 * 365日 = 10,950 TB(約11 PB)です。これは従来のKafkaにとって極端なケースであり、保持期間が短くなることがよくあります。コスト制約のため、従来のKafkaではより一般的な30日間の保持を想定しましょう。
- 総ストレージ: 30 TB/日 * 30日 = 900 TB(レプリケート済み)
- EBSボリュームタイプ:
gp3(費用対効果が高く、バランスの取れたパフォーマンス) - EBS
gp3コスト: $0.08/GB-月 - 総EBSコスト: 900,000 GB * 0.08/GB-月 = **72,000/月**
この計算はストレージのみをカバーしています。EC2インスタンスコスト、ネットワーク転送、IOPSは含まれていません。高いストレージコストは、多くの場合、組織に保持期間を大幅に短縮することを強制します。
階層型ストレージKafka
階層型ストレージを使用すると、ローカルに30日間、リモートに1年間データを保持できます。
-
ローカルストレージ (30日間):
- 総ストレージ: 900 TB(レプリケート済み)
- このストレージはブローカー全体に分散されます。各ブローカーが100GiBのローカル保持を持ち、10ブローカーがある場合、合計ローカルストレージは1TBです。これは簡略化されたもので、
local.retention.bytesはパーティションごとです。アクティブデータ用にブローカーあたり100GiBと仮定しましょう。 - 10ブローカー * 100 GiB/ブローカー = 1000 GiB = 1 TB
- EBS
gp3コスト: 1,000 GB * 0.08/GB-月 = **80/月**
-
リモートストレージ (1年間):
- 総生データ: 10 TB/日 * 365日 = 3,650 TB(非レプリケート、S3/GCSが耐久性を処理するため)
- AWS S3 Standardコスト: 最初の50TBは0.023/GB、次の450TBは0.022/GBなど(この規模では平均$0.022/GB-月)
- 総S3ストレージコスト: 3,650,000 GB * 0.022/GB-月 = **80,300/月**
- GCS Standardコスト: 最初の1PBは0.020/GBなど(この規模では平均0.020/GB-月)
- 総GCSストレージコスト: 3,650,000 GB * 0.020/GB-月 = **73,000/月**
-
オブジェクトストレージAPI呼び出し:
- 1GBセグメント、10TB/日のインジェストと仮定 = 10,000セグメント/日。
- S3 PUTリクエスト: 1,000リクエストあたり0.005。10,000 * 30日 = 300,000 PUT。コスト: 0.005 * 300 = $1.5/月。(無視できる)
- S3 GETリクエスト: 変動が大きい。コールドデータの10%が毎月アクセスされる場合(例: 365TB/月)、各GETが1GBセグメントの場合: 365,000 GET。コスト: 1,000リクエストあたり0.0004。0.0004 * 365 = $0.146/月。(無視できる)
- GCSの操作コストも同様で、多くの場合わずかに低い。
-
階層型ストレージの総コスト (AWS S3): 80 (ローカルEBS) + 80,300 (S3ストレージ) + ~2 (API呼び出し) = **~80,382/月**
-
階層型ストレージの総コスト (GCS): 80 (ローカルEBS) + 73,000 (GCSストレージ) + ~2 (API呼び出し) = **~73,082/月**
75%以上の節約の実証
この比較は、従来のKafkaでは1年間の保持ができないことが多いため、難しいです。EBSでの30日間の保持と、階層型ストレージでの1年間の保持を比較してみましょう。
- 従来型 (30日間 EBS): $72,000/月
- 階層型 (1年間 S3): $80,382/月
これでは節約は示されず、大幅に保持期間を延長したことによるコスト増を示しています。真の節約は、長期データのGB-月あたりのコストにあります。
シナリオを再評価しましょう: EBSで1年間のデータを保存しなければならなかったらどうなるでしょうか?
- 1年間のEBS総ストレージ: 10,950 TB(レプリケート済み)
- EBS
gp3コスト: 10,950,000 GB * 0.08/GB-月 = **876,000/月**
これを1年間の階層型ストレージと比較します。
- 階層型ストレージ (AWS S3): ~$80,382/月
- 階層型ストレージ (GCS): ~$73,082/月
節約計算 (AWS S3 vs. EBS 1年間保持): (876,000 - 80,382) / $876,000 = 0.908 = 90.8%の節約
節約計算 (GCS vs. EBS 1年間保持): (876,000 - 73,082) / $876,000 = 0.916 = 91.6%の節約
これは、長期保持の場合、階層型ストレージがオールEBSソリューションと比較してストレージコストを90%以上節約できることを示しています。EBSで3ヶ月の保持しか必要なかったとしても、コストは30日間のEBS + 1年間のS3/GCS階層型アプローチの3倍になります。
さらに、ブローカーのコンピューティングコストも削減できます。ローカルディスクの要件が小さくなるため、より小さなEC2インスタンスタイプを使用したり、ブローカーの数を減らしたりできる可能性があり、さらなる節約につながります。
パフォーマンス特性とトレードオフ
コスト削減は大きいですが、オブジェクトストレージ層を導入することによるパフォーマンスへの影響を理解することが重要です。
レイテンシとスループット
- ローカル読み取り:
local.retention期間内のデータの場合、パフォーマンスは従来のKafkaと同一です。レイテンシは通常1ミリ秒未満です。 - リモート読み取り: コンシューマーがオフロードされたデータを要求した場合、ブローカーはS3/GCSからフェッチする必要があります。これにより、ネットワークレイテンシとオブジェクトストレージの取得オーバーヘッドが発生します。
- 典型的なレイテンシ: ネットワーク状況、オブジェクトサイズ、クラウドプロバイダーのパフォーマンスに応じて、セグメントのフェッチあたり50ms〜500ms。これはローカルディスクアクセスよりも大幅に高くなります。
- 影響: 履歴データを読み取るコンシューマーは、より高いレイテンシと潜在的に低いスループットを経験します。これに敏感なアプリケーションは、ワーキングセットがローカル保持期間内にあることを確認する必要があります。
- 書き込み: オフロードは非同期かつ非ブロッキングであるため、プロデューサーの書き込みパフォーマンスはほとんど影響を受けません。
- ブローカーのCPU/ネットワーク: ブローカーは、リモートセグメントの管理に追加のCPUを消費し、オブジェクトストレージとの間でデータをアップロード/ダウンロードするためにネットワーク帯域幅を消費します。これはインスタンスサイジングに考慮する必要があります。
耐久性と可用性
- オブジェクトストレージの耐久性: AWS S3とGCSは、リージョン内の複数の施設にデータを冗長的に保存することで、極めて高い耐久性(11ナイン、99.999999999%)を提供します。これは一般的なEBS RAID構成よりも優れています。
- オブジェクトストレージの可用性: S3とGCSの両方で、高い可用性(標準ティアで99.9%から99.99%のSLA)を提供しています。
- Kafkaの耐久性: ローカルセグメントの3倍のレプリケーションファクターは引き続き適用されます。セグメントがオブジェクトストレージに正常にアップロードされると、その耐久性は実質的にクラウドプロバイダーによって処理されます。
比較表: KafkaストレージにおけるEBS vs. S3/GCS
| 機能 | 従来のKafka (EBS/Persistent Disk) | Kafka Tiered Storage (ローカルEBS + リモートS3/GCS)
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

実用的なGCPアーキテクチャガイド:実際に使うべきサービスと避けるべきもの
Google Cloud Platformの実践的な本番環境ガイド。Cloud RunがGKEを上回る理由、BigQueryとSecret Managerの活用方法、クラウド予算を食い荒らす5つの隠れたコストの罠を学びます。
Read more
BigQuery + Cloud Run: 本番向けのサーバーレスデータ取込パイプライン構築
Google Cloud 上でサーバーレスなデータ取込を本番品質で構築する実践ガイド。BigQuery Storage Write API、パーティショニングとクラスタリングの設計、Cloud Run 上の非同期 FastAPI レシーバ、Terraform による IaC 全体、実測に基づくコスト分析、そして深夜3時に呼ばれる障害モードまで扱います。
Read more
AI向けプラットフォームエンジニアリング: 自律エージェントのためのインフラストラクチャ設計
フリート規模のAIエージェントインフラストラクチャを設計するためのDevOpsガイド。OpenTelemetryトレーシング、Firecracker実行サンドボックス、ステートマシン、コストサーキットブレーカーについて解説。
Read more