Apache Kafka から Dataflow に読み込む

このドキュメントでは、Apache Kafka から Dataflow にデータを読み取る方法について説明します。また、パフォーマンスに関するヒントとベスト プラクティスについても説明します。

ほとんどのユースケースでは、マネージド I/O コネクタを使用して Kafka から読み取ることを検討してください。

より高度なパフォーマンスの調整が必要な場合は、KafkaIO コネクタの使用を検討してください。KafkaIO コネクタは、Java で利用できます。また、多言語パイプライン フレームワークを使用すると Python と Go でも利用できます。

並列処理

以降のセクションでは、Kafka から読み取るときに並列処理を構成する方法について説明します。

概要

並列処理は、ワーカーの最大数(max_num_workers)と Kafka パーティション数の 2 つの要因によって制限されます。Dataflow のデフォルトは、4 つの max_num_workers の並列処理ファンアウトです。ただし、ファンアウトはパーティションの数によって制限されます。たとえば、100 個の vCPU が使用可能であっても、パイプラインが 10 個の Kafka パーティションからのみ読み取る場合、最大並列処理は 10 になります。

並列処理を最大化するには、少なくとも 4 つの max_num_workers Kafka パーティションを用意することをおすすめします。ジョブで Runner v2 を使用している場合は、並列処理をさらに高く設定することを検討してください。ワーカー vCPU の数の 2 倍のパーティションを作成することをおすすめします。

再分配

パーティション数を増やせない場合は、KafkaIO.Read.withRedistribute を呼び出して並列処理を増やすことができます。この方法では、パイプラインに Redistribute 変換を追加することで、データをより効率的に再分散して並列化できるように Dataflow にヒントが提供されます。KafkaIO.Read.withRedistributeNumKeys を呼び出して、最適なシャード数を指定することを強くおすすめします。KafkaIO.Read.withRedistribute のみを単独で使用すると、多数のキーが生成され、パフォーマンスの問題が発生する可能性があります。詳細については、並列処理のレベルが高いステージを特定するをご覧ください。データを再分散すると、シャッフル ステップの実行で追加のオーバーヘッドが発生します。詳細については、融合を防ぐをご覧ください。

再分配シャッフルにかかる費用を最小限に抑えるには、KafkaIO.Read.withOffsetDeduplication を呼び出します。このモードでは、シャッフルの一部として永続化する必要があるデータ量を最小限に抑えつつ、exactly-once 処理を実現します。

exactly-once 処理が必要ない場合は、KafkaIO.Read.withAllowDuplicates を呼び出して重複を許可できます。

次の表に、再配布オプションの概要を示します。

オプション 処理モード Apache Beam 構成
入力の再配布 Exactly-once v2.60 以上 KafkaIO.Read.withRedistribute()
重複を許可 At-least-once v2.60 以上 KafkaIO.Read.withRedistribute().withAllowDuplicates()
オフセットの重複除去 Exactly-once v2.69 以上 KafkaIO.Read.withRedistribute().withOffsetDeduplication()

負荷の偏り

パーティション間の負荷が比較的均等で、偏らないようにします。負荷が偏っていると、ワーカーの使用率が低下する可能性があります。負荷の軽いパーティションから読み取るワーカーは比較的アイドル状態になる可能性がありますが、負荷の高いパーティションから読み取るワーカーは遅延する可能性があります。Dataflow は、パーティションごとのバックログの指標を提供します。

負荷が偏っている場合は、動的作業分散を使用して作業を分散できます。たとえば、Dataflow は、複数の少量パーティションから読み取るワーカーを 1 つ割り当て、単一の大量パーティションから読み取るワーカーを別のワーカーに割り当てることができます。ただし、2 つのワーカーが同じパーティションから読み取ることはできないため、負荷の高いパーティションは依然としてパイプラインの遅延の原因となります。

ベスト プラクティス

このセクションでは、Kafka から Dataflow に読み取る際の推奨事項について説明します。

少量のトピック

一般的なシナリオは、少量のトピック(お客様ごとに 1 つのトピックなど)を同時に読み取ることです。トピックごとに個別の Dataflow ジョブを作成することは、各ジョブに少なくとも 1 つの完全なワーカーが必要なため、費用対効果が低くなります。代わりに、次のオプションを検討してください。

  • トピックを統合する。トピックを Dataflow に取り込む前に結合します。少数の大量トピックを取り込む方が、多数の小量トピックを取り込むよりもはるかに効率的です。大量トピックごとに、ワーカーを最大限に活用する単一の Dataflow ジョブで処理できます。

  • 複数のトピックを読み取る。トピックを Dataflow に取り込む前に結合できない場合は、複数のトピックから読み取るパイプラインの作成を検討してください。このアプローチでは、Dataflow が複数のトピックを同じワーカーに割り当てることができます。アプローチを実施するには、次の 2 つの方法があります。

    • 単一の読み取りステップ。KafkaIO コネクタの単一インスタンスを作成し、複数のトピックを読み取るように構成します。次に、トピック名でフィルタして、トピックごとに異なるロジックを適用します。コード例については、複数のトピックから読み取るをご覧ください。すべてのトピックが同じクラスタに配置されている場合は、このオプションを検討してください。欠点としては、1 つのシンクまたは変換に問題があると、すべてのトピックがバックログに蓄積される可能性があります。

      高度なユースケースの場合は、読み取るトピックを指定する KafkaSourceDescriptor オブジェクトのセットを渡します。KafkaSourceDescriptor を使用すると、必要に応じてトピックリストを後で更新できます。この機能には、Runner v2 と Java が必要です。

    • 複数の読み取りステップ。異なるクラスタにあるトピックから読み取るには、パイプラインに複数の KafkaIO インスタンスを含めることができます。ジョブの実行中に、変換マッピングを使用して個々のソースを更新できます。新しいトピックまたはクラスタの設定は、Runner v2 を使用している場合にのみサポートされます。このアプローチでは、パイプライン レベルの指標に依存するのではなく、個々の読み取り変換を個別にモニタリングする必要があるため、オブザーバビリティが課題になる可能性があります。

Kafka への commit の戻し

デフォルトでは、KafkaIO コネクタは進行状況の追跡に Kafka オフセットを使用せず、Kafka に commit を返しません。commitOffsetsInFinalize を呼び出すと、コネクタは Dataflow でレコードが commit された後、Kafka に commit を返すよう試みます。Dataflow で commit されたレコードが完全に処理されていない場合があるため、パイプラインをキャンセルすると、レコードが完全に処理されずにオフセットが commit される可能性があります。

enable.auto.commit=True の設定は、Dataflow による処理なしで Kafka から読み取られた直後にオフセットを commit するため、このオプションの使用はおすすめしません。enable.auto.commit=False と commitOffsetsInFinalize=True の両方を設定することをおすすめします。enable.auto.commit を True に設定すると、処理中にパイプラインが中断された場合にデータが失われる可能性があります。すでに Kafka に commit されたレコードが破棄される可能性があります。

ウォーターマーク

デフォルトでは、KafkaIO コネクタは現在の処理時間を使用して、出力ウォーターマークとイベント時刻を割り当てます。この動作を変更するには、withTimestampPolicyFactory を呼び出して TimestampPolicy を割り当てます。Beam には、Kafka のログ追加時間またはメッセージの作成時間に基づいてウォーターマークを計算する TimestampPolicy が実装されています。

ランナーに関する考慮事項

KafkaIO コネクタには、Kafka 読み取り用の基盤となる実装が 2 つあります。古い ReadFromKafkaViaUnbounded と新しい ReadFromKafkaViaSDF です。Dataflow は、SDK 言語とジョブの要件に基づいて、ジョブに最適な実装を自動的に選択します。ランナーまたは Kafka の実装でのみ利用可能な特定の機能が必要な場合を除き、ランナーまたは Kafka の実装を明示的にリクエストしないでください。ランナーの選択の詳細については、Dataflow Runner v2 を使用するをご覧ください。

パイプラインで withTopic または withTopics を使用している場合、古い実装では、パイプラインの構築時に使用可能なパーティションについて Kafka にクエリを実行します。パイプラインを作成するマシンには、Kafka に接続する権限が必要です。権限エラーが表示された場合は、ローカルで Kafka に接続する権限があることを確認します。この問題を回避するには、withTopicPartitions を使用します。これは、パイプラインの構築時に Kafka に接続しません。

本番環境にデプロイする

ソリューションを本番環境にデプロイする場合は、Flex テンプレートを使用することをおすすめします。Flex テンプレートを使用すると、パイプラインが一貫した環境から起動されるため、ローカル構成の問題を軽減できます。

KafkaIO からのロギングは非常に冗長になる可能性があります。本番環境では、次のようにロギングレベルを下げることを検討してください。

sdkHarnessLogLevelOverrides='{"org.apache.kafka.clients.consumer.internals.SubscriptionState":"WARN"}'.

詳細については、パイプラインのワーカーログ レベルを設定するをご覧ください。

ネットワークを構成する

デフォルトでは、Dataflow はデフォルトの Virtual Private Cloud(VPC)ネットワーク内でインスタンスを起動します。Kafka の構成によっては、Dataflow に異なるネットワークとサブネットを構成しなければならない場合があります。詳細については、ネットワークとサブネットワークの指定をご覧ください。ネットワークを構成するときは、Dataflow ワーカーマシンが Kafka ブローカーにアクセスできるようにファイアウォール ルールを作成します。

VPC Service Controls を使用している場合は、Kafka クラスタを VPC Service Controls の境界内に配置するか、認可済み VPN または Cloud Interconnect に境界を拡張します。

Kafka クラスタが Google Cloudの外部にデプロイされている場合は、Dataflow と Kafka クラスタの間にネットワーク接続を作成する必要があります。複数のネットワーク オプションがあり、それぞれトレードオフがあります。

予測可能なパフォーマンスと信頼性という点で Dedicated Interconnect が最良のオプションですが、サードパーティが新しい回線をプロビジョニングする必要があるためセットアップに時間がかかることがあります。パブリック IP ベースのトポロジでは、必要となるネットワーキング作業がほとんどないため、すぐに使い始めることができます。

以降のセクションでは、これらのオプションについて詳しく説明します。

共有 RFC 1918 アドレス空間

Dedicated Interconnect と IPsec VPN の両方式は、Virtual Private Cloud(VPC)の RFC 1918 IP アドレスに直接アクセスできるため、Kafka の構成を簡素化できます。VPN ベースのトポロジを使用している場合は、ハイスループット VPN のセットアップを検討してください。

デフォルトでは、Dataflow はデフォルトの VPC ネットワークでインスタンスを起動します。プライベート ネットワーク トポロジで、Cloud Router で明示的に定義されたルートによって Google Cloud 内のサブネットワークが Kafka クラスタに接続されている場合は、Dataflow インスタンスをどこに配置するかを自分で制御できることが必要になります。Dataflow を使用して、network と subnetwork の実行パラメータを構成できます。

Dataflow がスケールアウトを試みてインスタンスを起動するとき、対応するサブネットワークで十分な数の IP アドレスを使用できることを確認してください。また、Dataflow インスタンスを起動するために別のネットワークを作成する場合は、プロジェクト内のすべての仮想マシン間の TCP トラフィックを有効にするファイアウォール ルールを設定するようにしてください。このファイアウォール ルールは、デフォルトのネットワークにはすでに構成されています。

パブリック IP アドレス空間

このアーキテクチャでは、Transport Layer Security(TLS)を使用して外部クライアントと Kafka との間のトラフィックを保護し、ブローカー間の通信には非暗号化トラフィックを使用します。Kafka リスナーが、内部通信と外部通信の両方に使用されるネットワーク インターフェースにバインドする場合、リスナーを構成するのは簡単です。ただし、多くのシナリオでは、外部にアドバタイズされる、クラスタ内の Kafka ブローカーのアドレスは、Kafka が使用する内部ネットワーク インターフェースのものとは異なります。このような場合は、advertised.listeners プロパティを使用できます。

# Configure protocol map
listener.security.protocol.map=INTERNAL:PLAINTEXT,EXTERNAL:SSL
# Use plaintext for inter-broker communication inter.broker.listener.name=INTERNAL
# Specify that Kafka listeners should bind to all local interfaces listeners=INTERNAL://0.0.0.0:9092,EXTERNAL://0.0.0.0:9093
# Separately, specify externally visible address advertised.listeners=INTERNAL://kafkabroker-n.mydomain.com:9092,EXTERNAL://kafkabroker-n.mydomain.com:9093

外部クライアントはポート 9093 を使用して SSL チャネル経由で接続し、内部クライアントはポート 9092 を使用して平文チャネル経由で接続します。advertised.listeners でアドレスを指定する際は、外部トラフィックと内部トラフィックの両方で同じインスタンスに解決される DNS 名を使用してください(このサンプルの場合は kafkabroker-n.mydomain.com)。パブリック IP アドレスは内部トラフィックでは解決できない可能性があるため、パブリック IP アドレスを使用すると機能しない場合があります。

Kafka をチューニングする

Kafka クラスタと Kafka クライアントの設定は、パフォーマンスに大きな影響を与える可能性があります。特に、次の設定が低すぎる可能性があります。このセクションでは、最初のステップとして推奨される値をいくつか説明しますが、特定のワークロードに合わせてこれらの値をテストする必要があります。

  • unboundedReaderMaxElements。デフォルトは 10,000 です。100,000 などの大きな値にすると、バンドルのサイズが大きくなる可能性があります。これによって、パイプラインに集計が含まれている場合、パフォーマンスが大幅に向上する可能性があります。ただし、値を大きくするとレイテンシも増加する可能性があります。値を設定するには、setUnboundedReaderMaxElements を使用します。この設定は Runner v2 には適用されません。Runner v2 の場合は、Dataflow サービス オプション sdf_checkpoint_after_output_bytes を使用します。

  • unboundedReaderMaxReadTimeMs。デフォルトは 10,000 ミリ秒です。20,000 ミリ秒などの大きな値にするとバンドルサイズが大きくなりますが、5,000 ミリ秒などの小さな値にするとレイテンシやバックログを削減できます。値を設定するには、setUnboundedReaderMaxReadTimeMs を使用します。この設定は Runner v2 には適用されません。Runner v2 の場合は、Dataflow サービス オプション sdf_checkpoint_after_duration を使用します。

  • max.poll.records。デフォルトは 500 です。値を大きくすると、特に Runner v2 を使用している場合、より多くの受信レコードをまとめて取得できるため、パフォーマンスが向上する可能性があります。値を設定するには、withConsumerConfigUpdates を呼び出します。

  • fetch.max.bytes。デフォルト値は 1 MB です。値を大きくすると、特に Runner v2 を使用する場合、リクエスト数が削減され、スループットが向上する可能性があります。ただし、設定値を高く設定しすぎるとレイテンシが増加する可能性があると同時に、ダウンストリーム処理が主なボトルネックになる可能性が高くなります。推奨される開始値は 100 MB です。値を設定するには、withConsumerConfigUpdates を呼び出します。

  • max.partition.fetch.bytes。デフォルト値は 1 MB です。このパラメータは、サーバーが返すパーティションあたりのデータの最大量を設定します。値を増やすと、特に Runner v2 を使用する場合、リクエスト数を削減してスループットを改善できます。ただし、設定値を高く設定しすぎるとレイテンシが増加する可能性があると同時に、ダウンストリーム処理が主なボトルネックになる可能性が高くなります。推奨される開始値は 100 MB です。値を設定するには、withConsumerConfigUpdates を呼び出します。

  • consumerPollingTimeout。デフォルトは 2 秒です。レコードを読み取る前にコンシューマ クライアントがタイムアウトする場合は、より大きい値を設定してみてください。この設定は、クロスリージョン読み取りやネットワーク速度の遅い読み取りを行う場合に最も関連性があります。値を設定するには、withConsumerPollingTimeout を呼び出します。

receive.buffer.bytes がメッセージのサイズを処理するのに十分な大きさであることを確認します。値が小さすぎると、コンシューマが継続的に再作成され、特定のオフセットを目標にしていることがログに表示されることがあります。

例

次のコードサンプルは、Kafka から読み取る Dataflow パイプラインの作成方法を示しています。アプリケーションのデフォルト認証情報を Google Cloud Managed Service for Apache Kafka 提供のコールバック ハンドラと組み合わせて使用する場合は、kafka-clients バージョン 3.7.0 以降が必要です。

単一のトピックから読み取る

この例では、マネージド I/O コネクタを使用します。Kafka トピックから読み取り、メッセージ ペイロードをテキスト ファイルに書き込む方法を示しています。

Java

Dataflow で認証を行うには、アプリケーションのデフォルト認証情報を設定します。詳細については、ローカル開発環境の認証を設定するをご覧ください。

import com.google.common.collect.ImmutableMap;
import java.io.UnsupportedEncodingException;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineResult;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.Description;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.StreamingOptions;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.values.TypeDescriptors;

public class KafkaRead {

  public static Pipeline createPipeline(Options options) {

    // Create configuration parameters for the Managed I/O transform.
    ImmutableMap<String, Object> config = ImmutableMap.<String, Object>builder()
        .put("bootstrap_servers", options.getBootstrapServer())
        .put("topic", options.getTopic())
        .put("format", "RAW")
        .put("max_read_time_seconds", 15)
        .put("auto_offset_reset_config", "earliest")
        .build();

    // Build the pipeline.
    var pipeline = Pipeline.create(options);
    pipeline
        // Read messages from Kafka.
        .apply(Managed.read(Managed.KAFKA).withConfig(config)).getSinglePCollection()
        // Get the payload of each message and convert to a string.
        .apply(MapElements
            .into(TypeDescriptors.strings())
            .via((row -> {
              var bytes = row.getBytes("payload");
              try {
                return new String(bytes, "UTF-8");
              } catch (UnsupportedEncodingException e) {
                throw new RuntimeException(e);
              }
            })))
        // Write the payload to a text file.
        .apply(TextIO
            .write()
            .to(options.getOutputPath())
            .withSuffix(".txt")
            .withNumShards(1));
    return pipeline;
  }
}

Python

Dataflow で認証を行うには、アプリケーションのデフォルト認証情報を設定します。詳細については、ローカル開発環境の認証を設定するをご覧ください。

import argparse

import apache_beam as beam

from apache_beam import window
from apache_beam.io.textio import WriteToText
from apache_beam.options.pipeline_options import PipelineOptions


def read_from_kafka() -> None:
    # Parse the pipeline options passed into the application. Example:
    #     --topic=$KAFKA_TOPIC --bootstrap_server=$BOOTSTRAP_SERVER
    #     --output=$CLOUD_STORAGE_BUCKET --streaming
    # For more information, see
    # https://beam.apache.org/documentation/programming-guide/#configuring-pipeline-options
    class MyOptions(PipelineOptions):
        @staticmethod
        def _add_argparse_args(parser: argparse.ArgumentParser) -> None:
            parser.add_argument("--topic")
            parser.add_argument("--bootstrap_server")
            parser.add_argument("--output")

    options = MyOptions()
    with beam.Pipeline(options=options) as pipeline:
        (
            pipeline
            # Read messages from an Apache Kafka topic.
            | beam.managed.Read(
                beam.managed.KAFKA,
                config={
                  "bootstrap_servers": options.bootstrap_server,
                  "topic": options.topic,
                  "data_format": "RAW",
                  "auto_offset_reset_config": "earliest",
                  # The max_read_time_seconds parameter is intended for testing.
                  # Avoid using this parameter in production.
                  "max_read_time_seconds": 5
                }
            )
            # Subdivide the output into fixed 5-second windows.
            | beam.WindowInto(window.FixedWindows(5))
            | WriteToText(
                file_path_prefix=options.output, file_name_suffix=".txt", num_shards=1
            )
        )

複数のトピックから読み取る

この例では、KafkaIO コネクタを使用します。複数の Kafka トピックから読み取り、トピックごとに個別のパイプライン ロジックを適用する方法を示します。

高度なユースケースでは、読み取るトピックのリストを更新できるように、一連の KafkaSourceDescriptor オブジェクトを動的に渡します。この方法には、Runner v2 と Java が必要です。

Java

Dataflow で認証を行うには、アプリケーションのデフォルト認証情報を設定します。詳細については、ローカル開発環境の認証を設定するをご覧ください。

import java.util.List;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineResult;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.io.kafka.KafkaIO;
import org.apache.beam.sdk.options.Description;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.options.StreamingOptions;
import org.apache.beam.sdk.transforms.Filter;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.values.TypeDescriptors;
import org.apache.kafka.common.serialization.LongDeserializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.joda.time.Duration;
import org.joda.time.Instant;

public class KafkaReadTopics {

  public static Pipeline createPipeline(Options options) {
    String topic1 = options.getTopic1();
    String topic2 = options.getTopic2();

    // Build the pipeline.
    var pipeline = Pipeline.create(options);
    var allTopics = pipeline
        .apply(KafkaIO.<Long, String>read()
            .withTopics(List.of(topic1, topic2))
            .withBootstrapServers(options.getBootstrapServer())
            .withKeyDeserializer(LongDeserializer.class)
            .withValueDeserializer(StringDeserializer.class)
            .withMaxReadTime(Duration.standardSeconds(10))
            .withStartReadTime(Instant.EPOCH)
        );

    // Create separate pipeline branches for each topic.
    // The first branch filters on topic1.
    allTopics
        .apply(Filter.by(record -> record.getTopic().equals(topic1)))
        .apply(MapElements
            .into(TypeDescriptors.strings())
            .via(record -> record.getKV().getValue()))
        .apply(TextIO.write()
            .to(topic1)
            .withSuffix(".txt")
            .withNumShards(1)
        );

    // The second branch filters on topic2.
    allTopics
        .apply(Filter.by(record -> record.getTopic().equals(topic2)))
        .apply(MapElements
            .into(TypeDescriptors.strings())
            .via(record -> record.getKV().getValue()))
        .apply(TextIO.write()
            .to(topic2)
            .withSuffix(".txt")
            .withNumShards(1)
        );
    return pipeline;
  }
}

Python

Dataflow で認証を行うには、アプリケーションのデフォルト認証情報を設定します。詳細については、ローカル開発環境の認証を設定するをご覧ください。

import argparse

import apache_beam as beam

from apache_beam.io.kafka import ReadFromKafka
from apache_beam.io.textio import WriteToText
from apache_beam.options.pipeline_options import PipelineOptions


def read_from_kafka() -> None:
    # Parse the pipeline options passed into the application. Example:
    #   --bootstrap_server=$BOOTSTRAP_SERVER --output=$STORAGE_BUCKET --streaming
    # For more information, see
    # https://beam.apache.org/documentation/programming-guide/#configuring-pipeline-options
    class MyOptions(PipelineOptions):
        @staticmethod
        def _add_argparse_args(parser: argparse.ArgumentParser) -> None:
            parser.add_argument('--bootstrap_server')
            parser.add_argument('--output')

    options = MyOptions()
    with beam.Pipeline(options=options) as pipeline:
        # Read from two Kafka topics.
        all_topics = pipeline | ReadFromKafka(consumer_config={
                "bootstrap.servers": options.bootstrap_server
            },
            topics=["topic1", "topic2"],
            with_metadata=True,
            max_num_records=10,
            start_read_time=0
        )

        # Filter messages from one topic into one branch of the pipeline.
        (all_topics
            | beam.Filter(lambda message: message.topic == 'topic1')
            | beam.Map(lambda message: message.value.decode('utf-8'))
            | "Write topic1" >> WriteToText(
                file_path_prefix=options.output + '/topic1/output',
                file_name_suffix='.txt',
                num_shards=1))

        # Filter messages from the other topic.
        (all_topics
            | beam.Filter(lambda message: message.topic == 'topic2')
            | beam.Map(lambda message: message.value.decode('utf-8'))
            | "Write topic2" >> WriteToText(
                file_path_prefix=options.output + '/topic2/output',
                file_name_suffix='.txt',
                num_shards=1))

次のステップ