•22 min read

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

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

Apache Kafkaの従来のアーキテクチャでは、コンピューティングとストレージが密接に結合されています。ブローカーはローカルのログセグメントを管理し、通常はAWS EBSやGCP Persistent Diskのような高性能ブロックストレージによってバックアップされています。この設計は、最近のデータに対しては比類のないスループットと低遅延アクセスを提供しますが、特にデータ量がペタバイト規模に拡大すると、長期的なデータ保持には経済的に持続不可能になります。履歴データやアクセス頻度の低いデータに対して、ブロックデバイスのプロビジョニングされたIOPSとストレージ容量のコストは、インフラ予算をすぐに圧迫する可能性があります。

KIP-405として導入されたKafka Tiered Storageは、Kafkaのストレージ層を根本的に再構築します。これにより、古くアクセス頻度の低いログセグメントを、ローカルのブローカーストレージから、AWS S3やGoogle Cloud Storage(GCS)のような費用対効果が高く、耐久性の高いオブジェクトストレージソリューションにシームレスにオフロードできます。このコンピューティングとストレージの分離により、Kafkaのコア保証を損なうことなく、大幅なコスト削減、スケーラビリティの向上、運用管理の簡素化が可能になります。

Audio Briefing
0:00 / 0:00

Kafka Tiered Storage (KIP-405) の理解

KIP-405は、2層のストレージモデルを導入しています。Kafkaブローカー上のローカル層と、オブジェクトストレージ上のリモート層です。主な目的は、ホットデータに対してローカルストレージのパフォーマンス上の利点を維持しつつ、履歴データに対してクラウドオブジェクトストレージの費用対効果と実質的に無限のスケーラビリティを活用することです。

コアコンセプト

  1. コンピューティングとストレージの分離: ブローカーは、すべてのトピックデータを保存する責任を単独で負うことはなくなります。主にローカルディスクから最新のデータを提供し、古いデータをリモート層にオフロードします。これにより、コンピューティング(ブローカー)とストレージ(オブジェクトストレージ)を独立してスケーリングできます。
  2. ローカルストレージ層: これは、ブローカーのローカルディスク上の従来のKafkaログディレクトリです。最新のログセグメントを保持し、アクティブなプロデューサーとコンシューマーに低遅延の読み書きアクセスを提供します。この層の保持期間は設定可能です。
  3. リモートストレージ層: これはクラウドオブジェクトストレージ(S3、GCS)です。設定された保持ポリシーに基づいてログセグメントが「コールド」と判断されると、非同期的にこの層にアップロードされます。この層は、すべての履歴データに対する不変で耐久性の高いアーカイブとして機能します。
  4. リモートストレージマネージャー: Kafkaブローカー内のコンポーネントで、ローカルストレージとリモートストレージ間のログセグメントのライフサイクルを管理します。セグメントのオブジェクトストレージへのアップロードと、コンシューマーがローカル保持期間を超えたデータを要求した際のセグメントの取得を処理します。
  5. メタデータ管理: 両方の層からのシームレスな消費を可能にするために、Kafkaはどのセグメントがローカルに存在し、どのセグメントがリモートストレージにオフロードされたかに関するメタデータを維持します。このメタデータは、コンシューマーが物理的な場所に関係なくデータを透過的にフェッチするために不可欠です。

階層型ストレージの利点

  • 大幅なコスト削減: オブジェクトストレージは、長期保持の場合、ブロックストレージ(EBS/Persistent Disk)よりも桁違いに安価です。これが導入の主な動機です。
  • スケーラビリティの向上: ストレージをブローカーから分離することで、ブローカーを追加したりローカルディスクサイズを増やしたりすることなく、ストレージ容量をほぼ無限にスケーリングできます。
  • 運用の簡素化: ブローカーは、より小さく高速なローカルディスクでプロビジョニングできるため、障害発生時の回復時間が短縮され、ディスク管理が簡素化されます。
  • データ保持期間の延長: 法外なコストをかけることなく、数ヶ月または数年間データを保持できるため、新しい分析ユースケースやコンプライアンス要件に対応できます。
  • ブローカーの回復性の向上: ローカルディスクのフットプリントが小さいほど、ブローカー障害発生時のデータレプリケーションと回復が高速になります。
Advertisement

アーキテクチャの詳細

階層型ストレージアーキテクチャは、Kafkaがログセグメントを管理する方法を根本的に変更します。データフローとコンポーネントの相互作用を理解することは、効果的なデプロイと最適化のために不可欠です。

データフローとコンポーネントの相互作用

  1. プロデューサーの書き込み: プロデューサーはメッセージをトピックパーティションに追加します。これらのメッセージは、ブローカーのローカルディスク上のアクティブなログセグメントに書き込まれます。
  2. セグメントのロールオーバー: アクティブなログセグメントがサイズ(log.segment.bytes)または時間(log.segment.ms)の制限に達すると、ロールオーバーされ、不変のクローズドセグメントになります。
  3. ローカル保持ポリシー: ブローカーは、ローカルに保持するセグメントを決定するためにlocal.retention.bytesまたはlocal.retention.msを適用します。このポリシーよりも古いセグメントは、オフロード対象としてマークされます。
  4. リモートセグメントのアップロード: RemoteLogManagerは、これらのマークされたセグメントを設定されたオブジェクトストレージバケット(S3またはGCS)に非同期的にアップロードします。このプロセスは、ブローカーの操作をブロックしません。
  5. ローカルセグメントの削除: セグメントがリモートストレージに正常にアップロードされ、そのメタデータが更新されると、ローカルコピーを削除してローカルディスクスペースを解放できます。
  6. コンシューマーの読み取り:
    • ホットデータ: コンシューマーがローカル保持期間内のデータを要求した場合、ブローカーはローカルディスクから直接データを提供し、低遅延を維持します。
    • コールドデータ: コンシューマーがリモートストレージにオフロードされたデータを要求した場合、ブローカーは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
  • 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のセットアップについて詳しく説明します。

前提条件

  1. Kafkaバージョン: KIP-405にはApache Kafka 3.0.0以降が必要です。本番環境では、安定性と機能のために3.3.x以降が推奨されます。
  2. クラウド認証情報:
    • 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ロールを持つサービスアカウントキー。
  3. クラウドバケット: 目的のリージョンに作成された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インスタンスタイプを使用したり、ブローカーの数を減らしたりできる可能性があり、さらなる節約につながります。

Advertisement

パフォーマンス特性とトレードオフ

コスト削減は大きいですが、オブジェクトストレージ層を導入することによるパフォーマンスへの影響を理解することが重要です。

レイテンシとスループット

  • ローカル読み取り: 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)

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement
BigQuery + Cloud Run: 本番向けのサーバーレスデータ取込パイプライン構築
gcp

BigQuery + Cloud Run: 本番向けのサーバーレスデータ取込パイプライン構築

Google Cloud 上でサーバーレスなデータ取込を本番品質で構築する実践ガイド。BigQuery Storage Write API、パーティショニングとクラスタリングの設計、Cloud Run 上の非同期 FastAPI レシーバ、Terraform による IaC 全体、実測に基づくコスト分析、そして深夜3時に呼ばれる障害モードまで扱います。

Read more