Pub/Sub ソースコネクタは、Pub/Sub から Kafka にメッセージをストリーミングします。これにより、Pub/Sub を Kafka ベースのアプリケーションやデータ パイプラインと統合できます。
Pub/Sub ソースコネクタのユースケースには以下が含まれます。
リアルタイムでのデータの取り込み。クラウド サービスや他のアプリケーションから Pub/Sub にデータをパブリッシュし、そのデータを Kafka に複製してストリーム処理を行います。
イベント ドリブン アーキテクチャ。Pub/Sub にパブリッシュされたメッセージから Kafka ベースの処理をトリガーします。
コネクタは、Pub/Sub サブスクリプションからメッセージを読み取り、各メッセージを Kafka レコードに変換して、Kafka トピックに書き込みます。デフォルトでは、コネクタは次のように Kafka レコードを作成します。
- Kafka レコードキーは
nullです。 - Kafka レコード値は、Pub/Sub メッセージ データ(バイト単位)です。
- Kafka レコード ヘッダーは空です。
ただし、この動作は構成できます。詳細については、 コネクタを構成するをご覧ください。
始める前に
Pub/Sub ソースコネクタを作成する前に、次のものがあることを確認してください。
サブスクリプションを含む Pub/Sub トピック。
Kafka トピックが Kafka クラスタ内にあります。
Connect クラスタ。 Connect クラスタを作成するときに、 Managed Service for Apache Kafka クラスタを プライマリ Kafka クラスタとして設定します。
必要なロールと権限
コネクタの作成に必要な権限を取得するには、プロジェクトに対するManaged Kafka コネクタ編集者 (roles/managedkafka.connectorEditor)IAM ロールの付与を管理者に依頼してください。ロールの付与については、プロジェクト、フォルダ、組織に対するアクセス権の管理をご覧ください。
この事前定義ロールには コネクタの作成に必要な権限が含まれています。必要とされる正確な権限については、必要な権限セクションを開いてご確認ください。
必要な権限
コネクタを作成するには、次の権限が必要です。
-
コネクタを作成する:
managedkafka.connectors.create
カスタムロールや他の事前定義ロールを使用して、これらの権限を取得することもできます。
Pub/Sub から読み取る権限を付与する
Managed Kafka サービス アカウントには、Pub/Sub サブスクリプションからメッセージを読み取る権限が必要です。Pub/Sub サブスクリプションを含むプロジェクトのサービス アカウントに次の IAM ロールを付与します。
- Pub/Sub サブスクライバー (
roles/pubsub.subscriber) - Pub/Sub 閲覧者 (
roles/pubsub.viewer)
Managed Kafka サービス アカウントの形式は
service-PROJECT_NUMBER@gcp-sa-managedkafka.iam.gserviceaccount.com です。ここで、PROJECT_NUMBER は Connect クラスタのプロジェクト番号です。
Connect クラスタが Managed Service for Apache Kafka クラスタとは異なるプロジェクトにある場合は、 別のプロジェクトに Connect クラスタを作成するをご覧ください。
Pub/Sub ソースコネクタを作成する
コンソール
コンソール Google Cloud で、[Connect クラスタ] ページに移動します。
コネクタを作成する Connect クラスタをクリックします。
[コネクタを作成] をクリックします。
コネクタ名に文字列を入力します。
コネクタの命名方法のガイドラインについては、 Managed Service for Apache Kafka リソースの命名ガイドラインをご覧ください。
[**コネクタ プラグイン**] で [**Pub/Sub ソース**] を選択します。
[Cloud Pub/Sub サブスクリプション] リストで、Pub/Sub サブスクリプションを選択します。コネクタは、このサブスクリプションからメッセージを pull します。サブスクリプションは、
projects/{project}/subscriptions/{subscription}という完全なリソース名で表示されます。[Kafka トピック] リストで、メッセージが書き込まれる Kafka トピックを選択します。
省略可: [構成] ボックスで、構成プロパティを追加するか、デフォルトのプロパティを編集します。詳細については、 コネクタを構成するをご覧ください。
[タスクの再起動ポリシー] を選択します。詳細については、 タスクの再起動ポリシーをご覧ください。
[作成] をクリックします。
gcloud
gcloud managed-kafka connectors createコマンドを実行します。gcloud managed-kafka connectors create CONNECTOR_ID \ --location=LOCATION \ --connect-cluster=CONNECT_CLUSTER_ID \ --config-file=CONFIG_FILE次のように置き換えます。
CONNECTOR_ID: コネクタの ID または名前。 コネクタの命名方法のガイドラインについては、 Managed Service for Apache Kafka リソースの命名ガイドラインをご覧ください。 コネクタの名前は変更できません。
LOCATION: Connect クラスタのロケーション。
CONNECT_CLUSTER_ID: コネクタが作成される Connect クラスタの ID。
CONFIG_FILE: YAML または JSON 構成ファイルへのパス。
以下に示すのは、構成ファイルの例です。
connector.class: "com.google.pubsub.kafka.source.CloudPubSubSourceConnector"
cps.project: "PROJECT_ID"
cps.subscription: "PUBSUB_SUBSCRIPTION_ID"
kafka.topic: "KAFKA_TOPIC_ID"
value.converter: "org.apache.kafka.connect.converters.ByteArrayConverter"
key.converter: "org.apache.kafka.connect.storage.StringConverter"
tasks.max: "3"
次のように置き換えます。
PROJECT_ID: Pub/Sub サブスクリプションが存在する Google Cloud プロジェクトの ID。
PUBSUB_SUBSCRIPTION_ID:データの pull 元となる Pub/Sub サブスクリプションの ID。
KAFKA_TOPIC_ID:データが書き込まれる Kafka トピックの ID。
cps.project、cps.subscription、kafka.topic
の構成プロパティが必要です。その他の構成オプションについては、
コネクタを構成するをご覧ください。
Terraform
Terraform リソースを使用して コネクタを作成できます。
Terraform 構成を適用または削除する方法については、 基本的な Terraform コマンドをご覧ください。
Go
このサンプルを試す前に、 クライアント ライブラリをインストールするにある Go の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Go API リファレンス ドキュメントをご覧ください。
Managed Service for Apache Kafka への認証を行うには、アプリケーションのデフォルト認証情報(ADC)を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。
Java
このサンプルを試す前に、 クライアント ライブラリをインストールするにある Java の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Java API リファレンス ドキュメントをご覧ください。
Managed Service for Apache Kafka への認証を行うには、アプリケーションのデフォルト認証情報を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。
Python
このサンプルを試す前に、 クライアント ライブラリをインストールするにある Python の設定手順を行ってください。詳細については、 Managed Service for Apache Kafka Python API リファレンス ドキュメントをご覧ください。
Managed Service for Apache Kafka への認証を行うには、アプリケーションのデフォルト認証情報を設定します。 詳細については、 ローカル開発環境の ADC の設定をご覧ください。
コネクタを作成した後は、コネクタの編集、削除、一時停止、停止、再起動を行うことができます。
コネクタを構成する
このセクションでは、コネクタに設定できる構成プロパティについて説明します。
このコネクタに固有のプロパティの一覧については、 Pub/Sub ソースコネクタの構成をご覧ください。
pull モード
pull モードでは、コネクタが Pub/Sub メッセージを取得する方法を指定します。次のモードがサポートされています。
pull モード (デフォルト)。メッセージはバッチで pull されます。このモードを有効にするには、
cps.streamingPull.enabled=false.を設定します。バッチサイズを構成するには、cps.maxBatchSizeプロパティを設定します。pull モードの詳細については、 Pull API をご覧ください。
Streaming Pull モード 。Pub/Sub からメッセージを取得する際の最大スループットと最小レイテンシを実現します。このモードを有効にするには、
cps.streamingPull.enabled=trueを設定します。Streaming Pull モードの詳細については、 StreamingPull API をご覧ください。
Streaming Pull が有効になっている場合は、次の構成プロパティを設定してパフォーマンスを調整できます。
cps.streamingPull.flowControlBytes: タスクあたりの未処理メッセージ バイト数の最大値。cps.streamingPull.flowControlMessages: タスクあたりの未処理メッセージ数の最大値。cps.streamingPull.maxAckExtensionMs: コネクタがサブスクライブの期限を延長する最大時間(ミリ秒単位)。cps.streamingPull.maxMsPerAckExtension: コネクタが延長ごとにサブスクライブの期限を延長する最大時間(ミリ秒単位)。cps.streamingPull.parallelStreams: サブスクリプションからメッセージを pull するストリームの数。
Pub/Sub エンドポイント
デフォルトでは、コネクタはグローバル Pub/Sub エンドポイントを使用します。エンドポイントを指定するには、cps.endpoint プロパティをエンドポイント アドレスに設定します。
エンドポイントの詳細については、
Pub/Sub エンドポイントをご覧ください。
Kafka パーティション
デフォルトでは、コネクタはトピック内の単一のパーティションに書き込みます。コネクタが書き込むパーティションの数を指定するには、kafka.partition.count プロパティを設定します。値は、トピックの
パーティション数を超えないようにしてください。
コネクタがメッセージをパーティションに割り当てる方法を指定するには、kafka.partition.scheme
プロパティを設定します。詳細については、
Pub/Sub ソースコネクタの構成をご覧ください。
コンバータ
キー コンバータを org.apache.kafka.connect.storage.StringConverter に設定します。
コネクタの構成に応じて、値コンバータを次のいずれかに設定します。
org.apache.kafka.connect.converters.ByteArrayConverterorg.apache.kafka.connect.json.JsonConverter
詳細については、レコード値をご覧ください。
メッセージ コンバージョン
Pub/Sub ソースコネクタは、Pub/Sub メッセージを Kafka レコードに変換します。以降のセクションでは、変換プロセスについて説明します。
レコードキー
キー コンバータは org.apache.kafka.connect.storage.StringConverter である必要があります。
デフォルトでは、レコードキーは
nullです。Pub/Sub メッセージ属性をキーとして使用するには、
kafka.key.attributeを属性の名前に設定します。例:kafka.key.attribute=usernamePub/Sub 順序指定キーをキーとして使用するには、
kafka.key.attribute=orderingKeyを設定します。
レコード ヘッダー
デフォルトでは、レコード ヘッダーは空です。
kafka.record.headers が true の場合、Pub/Sub メッセージ属性はレコード
ヘッダーとして書き込まれます。順序指定キーを含めるには、cps.makeOrderingKeyAttribute=true
を設定します。
レコード値
レコード値は、バイト配列または struct 型として書き込まれます。
バイト配列レコード値
kafka.record.headers が true の場合、または Pub/Sub
メッセージにカスタム属性がない場合、コネクタはメッセージ データをバイト配列として書き込みます。値コンバータを
org.apache.kafka.connect.converters.ByteArrayConverter に設定します。
構造体レコード値
kafka.record.headers が false で、メッセージに 1 つ以上のカスタム属性がある場合、コネクタはレコード値を
struct として書き込みます。値コンバータを org.apache.kafka.connect.json.JsonConverter
に設定します。
struct には次のフィールドが含まれます。
message: Pub/Sub メッセージ データ(バイト単位)。Pub/Sub メッセージ属性ごとのフィールド。順序指定キーを含めるには、
cps.makeOrderingKeyAttribute=trueを設定します。
たとえば、メッセージに username 属性がある場合、レコード値は次のようになります。
{
"message":"MESSAGE_DATA",
"username":"Alice"
}
value.converter.schemas.enable が true の場合、struct にはペイロードとスキーマの両方が含まれます。
{
"schema":
{
"type":"struct",
"fields": [
{
"type":"bytes",
"optional":false,
"field":"message"
},
{
"type":"string",
"optional":false,
"field":"username"
}
],
"optional":false
},
"payload": {
"message":"MESSAGE_DATA",
"username":"Alice"
}
}