MaxCompute sink コネクタを作成して、インスタンスのソーストピックから MaxCompute のテーブルにデータをエクスポートします。
前提条件
MaxCompute sink コネクタを作成する前に、両方のサービスで以下の準備を完了してください。
ApsaraMQ for Kafka — ソースデータを提供するインスタンスを準備します。
ApsaraMQ for Kafka インスタンスのコネクタ機能を有効にします。 詳細については、「コネクタの有効化」をご参照ください。
ApsaraMQ for Kafka インスタンスのソーストピックを作成します。 詳細については、「ステップ 1: トピックの作成」をご参照ください。
このトピックでは、`maxcompute-test-input` という名前のトピックを例として使用します。
MaxCompute — データを受信する宛先テーブルを準備します。
MaxCompute クライアントを使用してテーブルを作成します。 詳細については、「テーブルの作成」をご参照ください。
このトピックでは、`connector_test` という名前のプロジェクトの `test_kafka` という名前のテーブルを例として使用します。 次のステートメントでテーブルを作成します。
CREATE TABLE IF NOT EXISTS test_kafka(topic STRING,partition BIGINT,offset BIGINT,key STRING,value STRING) PARTITIONED by (pt STRING);オプション:EventBridge — EventBridge デプロイパスを使用するインスタンスにのみ必要です。
説明コネクタタスクが属するインスタンスが中国 (杭州) または中国 (成都) リージョンにある場合にのみ、この操作を完了する必要があります。
注意事項
MaxCompute sink コネクタを作成する前に、以下の制限と動作を確認してください。
タスクがサポートされるかどうかを決定する制限
リージョン — ApsaraMQ for Kafka インスタンスのソーストピックから MaxCompute へのデータのエクスポートは、同一リージョン内でのみ可能です。 コネクタの制限事項の詳細については、「制限事項」をご参照ください。
インスタンスのメジャーバージョン 0.10.2 — MaxCompute sink コネクタで必要なトピックの一部は、ローカルストレージエンジンを使用する必要があります。 メジャーバージョンが 0.10.2 の ApsaraMQ for Kafka インスタンスは、ローカルストレージを使用するトピックの手動作成をサポートしていません。 これらのトピックは自動的にのみ作成できます。 ご利用のインスタンスがこのバージョンを実行している場合は、コネクタに必要なトピックを自動的に作成させてください。
中国 (杭州) および中国 (成都) リージョンのインスタンス
コネクタが属するインスタンスが中国 (杭州) または中国 (成都) リージョンにある場合、この機能は EventBridge にデプロイされます。 他のすべてのリージョンのインスタンスはデフォルトのデプロイパスを使用するため、以下の項目は適用されません。
課金 — EventBridge は現在無料です。 詳細については、「課金」をご参照ください。
サービスリンクロール — コネクタを作成すると、EventBridge は `AliyunServiceRoleForEventBridgeSourceKafka` と `AliyunServiceRoleForEventBridgeConnectVPC` のサービスリンクロールを自動的に作成します。 まだ存在しない各ロールに対して、作成ウィザードに **[サービス認証]** ダイアログボックスが表示されます。
サービスリンクロールが作成されていない場合、EventBridge は対応するサービスリンクロールを自動的に作成し、EventBridge がそのロールを使用して ApsaraMQ for Kafka および Virtual Private Cloud (VPC) にアクセスできるようにします。
サービスリンクロールがすでに作成されている場合、EventBridge はそれらを再度作成しません。
サービスリンクロールの詳細については、「EventBridge のサービスリンクロール」をご参照ください。
タスク実行ログ — EventBridge にデプロイされたタスクは、タスク実行ログの表示をサポートしていません。 コネクタタスクが完了した後、ソーストピックをサブスクライブする グループの消費詳細に基づいてタスクの進行状況を確認できます。 詳細については、「コンシューマーのステータスの表示」をご参照ください。
操作手順
MaxCompute sink コネクタを使用して ApsaraMQ for Kafka インスタンスのソーストピックから MaxCompute のテーブルにデータをエクスポートするには、以下のステージを順番に完了します。 各ステージは、以下のセクションのいずれかで説明されています。
ApsaraMQ for Kafka に MaxCompute へのアクセス権限を付与します。
オプション: MaxCompute sink コネクタで必要なトピックと グループを作成します。
トピックと グループをカスタマイズする必要がない場合は、このステージをスキップして、次のステージで自動作成を選択できます。 インスタンスのメジャーバージョンが 0.10.2 の場合は、このトピックの「注意事項」セクションで説明されているように、自動作成を使用します。
結果の検証
RAM ロールの作成
Resource Access Management (RAM) ロールは、信頼できるサービスとして ApsaraMQ for Kafka を直接選択することをサポートしていません。 したがって、RAM ロールを作成するときは、サポートされている任意のサービスを信頼できるサービスとして選択します。 RAM ロールが作成された後、手動で信頼ポリシーを変更します。
RAM コンソールにログインします。
左側のナビゲーションバーで、[ID] > [ロール] を選択します。
[ロール] ページで、[ロールの作成] をクリックします。
次の図は、[ロール] ページの [ロールの作成] ボタンを示しています。

[ロールの作成] パネルで、次の操作を実行します。
信頼できるエンティティタイプを [Alibaba Cloud サービス] に設定し、次へ をクリックします。
[ロールタイプ] セクションで、[通常のサービスロール] を選択します。 [ロール名] フィールドに `AliyunKafkaMaxComputeUser1` と入力します。 [信頼できるサービスの選択] リストから [MaxCompute] を選択し、[完了] をクリックします。
[ロール] ページで、[AliyunKafkaMaxComputeUser1] を見つけてクリックします。
[AliyunKafkaMaxComputeUser1] ページで、[信頼ポリシー] タブをクリックし、[信頼ポリシーの編集] をクリックします。
[信頼ポリシーの編集] パネルで、スクリプト内の odps を `alikafka` に置き換えてから、[OK] をクリックします。
次のポリシーは、置換後の結果を示しています。

権限の追加
コネクタがメッセージを MaxCompute テーブルに同期できるようにするには、作成した RAM ロールに少なくとも次の権限を付与する必要があります。
オブジェクト | アクション | 記述 |
プロジェクト | CreateInstance | プロジェクトにインスタンスを作成します。 |
テーブル | Describe | テーブルのメタデータを読み取ります。 |
テーブル | Alter | テーブルのメタデータを変更するか、パーティションを追加または削除します。 |
テーブル | Update | テーブルのデータを上書きまたは追加します。 |
上記の権限とその付与方法の詳細については、「MaxCompute の権限」をご参照ください。
以下の手順では、この Topic で作成された AliyunKafkaMaxComputeUser1 ロールに権限を付与する方法を説明します。 MaxCompute は RAM ロールをユーザーとして管理するため、このセクションのすべてのコマンドは RAM$<accountid>:role/aliyunkafkamaxcomputeuser1 オブジェクトを対象とします。
以下の各コマンドで、<accountid> をご自身の Alibaba Cloud アカウント ID に置き換えてください。
MaxCompute クライアントにログインします。
次のコマンドを実行して、RAM ロールをユーザーとして追加します。
add user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;RAM ロールに MaxCompute へのアクセスに必要な最小権限を付与します。
次のコマンドを実行して、RAM ロールにプロジェクトに対する権限を付与します。
grant CreateInstance on project connector_test to user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;次のコマンドを実行して、RAM ロールにテーブルに対する権限を付与します。
grant Describe, Alter, Update on table test_kafka to user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;MaxCompute sink コネクタで必要なトピックの作成
ApsaraMQ for Kafka コンソールでは、MaxCompute sink コネクタで必要な 5 つのトピック (タスクオフセットトピック、タスク設定トピック、タスクステータストピック、デッドレターキュートピック、エラーデータトピック) を手動で作成できます。 必要なパーティション数とストレージエンジンは、次の表に示すようにトピックによって異なります。
トピック | 推奨される名前のプレフィックス | パーティション | ストレージエンジン | cleanup.policy |
タスクオフセットトピック | connect-offset | 1 より大きい | ローカルストレージ | compact |
タスク設定トピック | connect-config | 1 | ローカルストレージ | compact |
タスクステータストピック | connect-status | 6 を推奨 | ローカルストレージ | compact |
デッドレターキュートピック | connect-error | 6 を推奨 | ローカルストレージまたはクラウドストレージ | — |
エラーデータトピック | connect-error | 6 を推奨 | ローカルストレージまたはクラウドストレージ | — |
トピックリソースを節約するために、1 つのトピックをデッドレターキュートピックとエラーデータトピックの両方として使用できます。 デッドレターキュートピックとエラーデータトピックの場合、ストレージエンジンはローカルストレージまたはクラウドストレージにすることができ、残りのプロパティは、このセクションのトピックプロパティテーブルで説明されている一般的なルールに従います。 各トピックの完全な説明については、「ソースサービスの設定タブのパラメーター」をご参照ください。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
概要 ページで、リソースの分布 セクションのリージョンを選択します。
注:Topic は、ECS インスタンスがデプロイされているリージョン、つまりご利用のアプリケーションと同じリージョンに作成する必要があります。Topic はリージョンをまたいで使用することはできません。例えば、Topic が中国 (北京) リージョンで作成された場合、メッセージプロデューサーとコンシューマーも、中国 (北京) リージョンの ECS インスタンス上で実行する必要があります。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
Topic は、ご利用のアプリケーションと同じリージョン、つまり Elastic Compute Service (ECS) インスタンスがデプロイされているリージョンに作成する必要があります。Topic はリージョンをまたいで使用することはできません。例えば、Topic が中国 (北京) リージョンで作成された場合、メッセージのプロデューサーとコンシューマーも、中国 (北京) リージョンの ECS インスタンス上で実行する必要があります。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
Topic は、アプリケーションと同じリージョン、つまり ECS インスタンスがデプロイされているリージョンに作成する必要があります。Topic はリージョンをまたいで使用することはできません。例えば、Topic が中国 (北京) リージョンで作成された場合、メッセージプロデューサーとコンシューマーも中国 (北京) リージョンの ECS インスタンス上で実行する必要があります。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
Topic は、アプリケーションと同じリージョン、つまり ECS インスタンスがデプロイされているリージョンに作成する必要があります。Topic は、リージョンをまたいで使用することはできません。例えば、Topic が中国 (北京) リージョンに作成された場合、メッセージのプロデューサーとコンシューマーも、中国 (北京) リージョンにある ECS インスタンス上で実行する必要があります。
インスタンスリスト ページで、ターゲットインスタンスの名前をクリックします。
インスタンスリスト ページで、ターゲットインスタンスの名前をクリックします。
インスタンスリスト ページで、ターゲットインスタンスの名前をクリックします。
左側のナビゲーションウィンドウで、トピック管理 をクリックします。
左側のナビゲーションウィンドウで、トピック管理 をクリックします。
トピック管理 ページで、トピックの作成 をクリックします。
各コネクタトピックのプロパティを前の表に記載されている値に設定します。 次の表に、すべてのトピックプロパティを示します。
パラメーター | 説明 | 例 |
名前 | Topic の名前です。Kafka では、 | connect-offset-kafka-maxcompute-sink |
記述 | トピックの簡単な説明。 | demo test |
パーティションの数 | トピックのパーティション数。 | 12 |
ストレージエンジン > 注: サーバーレスでない Professional Edition インスタンスのみがストレージエンジンタイプの選択をサポートします。 他のインスタンスはこのオプションをサポートせず、デフォルトでクラウドストレージタイプを使用します。 | トピックメッセージのストレージエンジン。 ApsaraMQ for Kafka は、次の 2 つのストレージエンジンをサポートしています。 - クラウドストレージ:基盤レイヤーで Alibaba Cloud ディスクを使用し、低レイテンシー、高性能、耐久性、高信頼性を提供します。 このエンジンは、分散 3 レプリカメカニズムを使用します。 インスタンスの 仕様タイプが Standard Edition (High Write) の場合、ストレージエンジンは クラウドストレージのみになります。 - ローカルストレージ:オープンソース Kafka の In-Sync Replicas (ISR) レプリケーションアルゴリズムと分散 3 レプリカメカニズムを使用します。 | ローカルストレージ |
メッセージタイプ | トピックメッセージのタイプ。 - 通常のメッセージ:デフォルトでは、同じキーを持つメッセージは同じパーティションに配布され、パーティション内のメッセージは送信された順序で保存されます。 クラスター内のマシンに障害が発生した場合、メッセージの順序が乱れる可能性があります。 ストレージエンジンを クラウドストレージに設定した場合、デフォルトで 通常のメッセージが選択されます。 - パーティション順位メッセージ:デフォルトでは、同じキーを持つメッセージは同じパーティションに配布され、パーティション内のメッセージは送信された順序で保存されます。 クラスター内のマシンに障害が発生した場合でも、パーティション内のメッセージは送信された順序で保存されます。 ただし、これらのパーティションが回復するまで、一部のパーティションへのメッセージの送信が失敗する可能性があります。 ストレージエンジンを ローカルストレージに設定した場合、デフォルトで パーティション順位メッセージが選択されます。 | 通常のメッセージ |
ログリリースポリシー | トピックログのクリーンアップポリシー。 ストレージエンジンを ローカルストレージに設定した場合、ログリリースポリシーを設定する必要があります。 ApsaraMQ for Kafka は、次の 2 つのクリーンアップポリシーをサポートしています。 - Delete:デフォルトのメッセージクリーンアップポリシー。 ディスク容量が十分な場合、メッセージは最大保持期間内に保持されます。 ディスク容量が不十分な場合 (通常はディスク使用率が 85% を超える場合)、サービス可用性を確保するために古いメッセージが予定より早く削除されます。 - Compact:Kafka ログコンパクションポリシーを使用します。 ログコンパクションポリシーは、各メッセージキーの最新の値が保持されることを保証します。 このポリシーは、システム障害後の状態復元やシステム再起動後のキャッシュ再読み込みなどのシナリオに適しています。 たとえば、Confluent Schema Registry または Kafka Connect を使用する場合、システムの状態または設定情報を保存するために Kafka コンパクション済みトピックを使用する必要があります。 > 重要: コンパクション済みトピックは、通常、Confluent Schema Registry や Kafka Connect などの特定の生態系コンポーネントにのみ使用されます。 他のシナリオでメッセージを送受信するために使用されるトピックには、このプロパティを設定しないでください。 詳細については、「ApsaraMQ for Kafka デモライブラリ」をご参照ください。 | Compact |
タグ | トピックのタグ。 | demo |
Topic が作成されると、トピック管理 ページのリストに表示されます。これらの Topic は手動で作成されているため、コネクタを作成する際は リソースの作成方法 を 手動で作成します に設定し、作成した Topic の名前を入力してください。
MaxCompute sink コネクタで必要な グループの作成
ApsaraMQ for Kafka コンソールでは、MaxCompute sink コネクタのデータ同期タスクで使用される グループを手動で作成できます。 グループの名前は `connect-タスク名` である必要があります。ここで、`タスク名` はコネクタの名前です。 詳細については、「ソースサービスの設定タブのパラメーター」をご参照ください。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
インスタンスリスト ページで、ターゲットインスタンスの名前をクリックします。
インスタンスリストページで、ターゲットインスタンスの名前をクリックします。
インスタンスリスト ページで、ターゲットインスタンスの名前をクリックします。
左側のナビゲーションウィンドウで、Group の管理 をクリックします。
左側のナビゲーションウィンドウで、Group の管理 をクリックします。
Group の管理 ページで、グループの作成 をクリックします。
グループが作成されると、Group の管理 ページのリストに表示されます。コネクタを作成する際に、Connector コンシューマーグループ フィールドにその名前を入力します。
MaxCompute sink コネクタの作成とデプロイ
ApsaraMQ for Kafka から MaxCompute にデータを同期する MaxCompute sink コネクタを作成してデプロイします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
ApsaraMQ for Kafka コンソールにログインします。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
概要ページのリソースの分布セクションで、リージョンを選択します。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
概要 ページの リソースの分布 セクションで、リージョンを選択します。
左側のナビゲーションウィンドウで、Connector タスクリスト をクリックします。
左側のナビゲーションウィンドウで、Connector タスクリスト をクリックします。
左側のナビゲーションウィンドウで、Connector タスクリスト をクリックします。
Connector タスクリスト ページで、インスタンスの選択 ドロップダウンリストから、コネクタが属するインスタンスを選択し、Connector の作成 をクリックします。
Connector タスクリスト ページで、インスタンスの選択 ドロップダウンリストからコネクタが属するインスタンスを選択し、Connector の作成をクリックします。
1. Connector の作成 構成ウィザードページで、以下の操作を実行します。
基本情報の設定 タブで、ビジネス要件に基づいて次のパラメーターを設定し、次へ をクリックします。
パラメーター
説明
例
名前
コネクタの名前。 命名規則: - 名前には数字、小文字、ハイフン (-) を含めることができますが、ハイフン (-) で始めることはできません。 名前の長さは 1~48 文字である必要があります。 - 名前は、同じ ApsaraMQ for Kafka インスタンス内で一意である必要があります。 コネクタのデータ同期タスクは、`connect-タスク名` という名前の グループを使用する必要があります。 グループを手動で作成しない場合、システムが自動的に作成します。
kafka-maxcompute-sink
インスタンス
デフォルトでは、インスタンス名とインスタンス ID が指定されます。
demo alikafka_post-cn-st21p8vj****
[ソースサービスの設定] タブで、[データソース] を [ApsaraMQ for Kafka] に設定し、次のパラメーターを設定してから 次へ をクリックします。
説明すでに Topic とコンシューマーグループを作成している場合は、手動でのリソース作成を選択し、既存リソースの情報を入力してください。それ以外の場合は、自動でのリソース作成を選択してください。
すでに Topic と コンシューマーグループを作成している場合は、手動リソース作成を選択し、既存リソースの情報を入力してください。そうでない場合は、自動リソース作成を選択してください。
すでに Topic と コンシューマーグループ を作成済みの場合は、手動でのリソース作成を選択し、既存リソースの情報を入力してください。それ以外の場合は、自動でのリソース作成を選択してください。
すでに Topic と コンシューマーグループ を作成している場合は、手動リソース作成を選択して、既存リソースの情報を入力します。それ以外の場合は、自動リソース作成を選択します。
すでに Topic と コンシューマーグループを作成している場合は、手動リソース作成を選択し、既存リソースの情報を入力してください。それ以外の場合は、自動リソース作成を選択してください。
パラメーター
説明
例
データソース Topic
同期したいデータのソーストピック。
maxcompute-test-input
コンシューマースレッドの同時発生数
ソーストピックの同時コンシューマースレッド数。 デフォルト値:6。 有効な値: - 1 - 2 - 3 - 6 - 12
6
消費の開始位置
メッセージが消費されるオフセット。 有効な値: - 一番古いオフセット:最小オフセットから消費を開始します。 - 一番新しいオフセット:最大オフセットから消費を開始します。
一番古いオフセット
VPC ID
データ同期タスクが実行される VPC。 このパラメーターを表示するには、実行環境の設定 をクリックします。 デフォルトでは、ApsaraMQ for Kafka インスタンスがデプロイされている VPC が使用されます。 このパラメーターを指定する必要はありません。
vpc-bp1xpdnd3l***
VSwitch ID
データ同期タスクが実行される vSwitch。 このパラメーターを表示するには、実行環境の設定 をクリックします。 vSwitch は、ApsaraMQ for Kafka インスタンスと同じ VPC にある必要があります。 デフォルトでは、ApsaraMQ for Kafka インスタンスをデプロイしたときに指定した vSwitch が使用されます。
vsw-bp1d2jgg81***
失敗の処理
メッセージの送信に失敗した後、エラーが発生したトピックのパーティションへのサブスクリプションを続行するかどうかを指定します。 このパラメーターを表示するには、実行環境の設定 をクリックします。 有効な値: - サブスクリプションの継続:エラーが発生したトピックのパーティションへのサブスクリプションを続行し、エラーログを出力します。 - サブスクリプションの停止:エラーが発生したトピックのパーティションへのサブスクリプションを停止し、エラーログを出力します > 注: - ログの表示方法の詳細については、「コネクタ操作」をご参照ください。 - エラーコードに基づいた解決策の検索方法の詳細については、「エラーコード」をご参照ください。 - エラーが発生したトピックのパーティションへのサブスクリプションを再開するには、して テクニカルサポートにお問い合わせください。
サブスクリプションの継続
リソースの作成方法
コネクタで必要なトピックと グループを作成する方法。 このパラメーターを表示するには、実行環境の設定 をクリックします。 有効な値: - 自動作成:システムが必要なトピックと グループを自動的に作成します。 事前に作成していない場合は、この値を選択します。 - 手動で作成します:事前にトピックと グループを作成した場合は、この値を選択し、次のフィールドにそれらの名前を入力します。
自動作成
Connector コンシューマーグループ
コネクタのデータ同期タスクで使用される グループ。 このパラメーターを表示するには、実行環境の設定 をクリックします。 グループの名前は `connect-タスク名` である必要があります。
connect-kafka-maxcompute-sink
タスクサイトの Topic
コンシューマーオフセットを保存するトピック。 このパラメーターを表示するには、実行環境の設定 をクリックします。 - トピック名:名前は connect-offset で始まることを推奨します。 - パーティション数:トピックのパーティション数は 1 より大きい必要があります。 - ストレージエンジン:トピックのストレージエンジンはローカルストレージである必要があります。 - cleanup.policy:トピックのログクリーンアップポリシーは compact である必要があります。
connect-offset-kafka-maxcompute-sink
タスク設定の Topic
タスク設定を保存するトピック。 このパラメーターを表示するには、実行環境の設定 をクリックします。 - トピック名:名前は connect-config で始まることを推奨します。 - パーティション数:トピックのパーティション数は 1 である必要があります。 - ストレージエンジン:トピックのストレージエンジンはローカルストレージである必要があります。 - cleanup.policy:トピックのログクリーンアップポリシーは compact である必要があります。
connect-config-kafka-maxcompute-sink
タスクステータスの Topic
タスクステータスを保存するトピック。 このパラメーターを表示するには、実行環境の設定 をクリックします。 - トピック名:名前は connect-status で始まることを推奨します。 - パーティション数:トピックのパーティション数を 6 に設定することを推奨します。 - ストレージエンジン:トピックのストレージエンジンはローカルストレージである必要があります。 - cleanup.policy:トピックのログクリーンアップポリシーは compact である必要があります。
connect-status-kafka-maxcompute-sink
デッドレターキューの Topic
Connect フレームワークの異常データを保存するトピック。 このパラメーターを表示するには、実行環境の設定 をクリックします。 トピックリソースを節約するために、このトピックをエラーデータトピックとして使用できます。 - トピック名:名前は connect-error で始まることを推奨します。 - パーティション数:トピックのパーティション数を 6 に設定することを推奨します。 - ストレージエンジン:トピックのストレージエンジンはローカルストレージまたはクラウドストレージにすることができます。
connect-error-kafka-maxcompute-sink
例外データ Topic
sink の異常データを保存するトピック。 このパラメーターを表示するには、実行環境の設定 をクリックします。 トピックリソースを節約するために、このトピックをデッドレターキュートピックとして使用できます。 - トピック名:名前は connect-error で始まることを推奨します。 - パーティション数:トピックのパーティション数を 6 に設定することを推奨します。 - ストレージエンジン:トピックのストレージエンジンはローカルストレージまたはクラウドストレージにすることができます。
connect-error-kafka-maxcompute-sink
ターゲットサービスの設定 タブで、宛先サービスを MaxCompute に設定し、以下のパラメーターを設定した後、作成 をクリックします。
コネクタが属するインスタンスが中国 (杭州) または中国 (成都) リージョンにある場合、[宛先サービス] を [MaxCompute] に設定すると、`AliyunServiceRoleForEventBridgeSourceKafka` および `AliyunServiceRoleForEventBridgeConnectVPC` サービスリンクロールのそれぞれに対して [サービス認証] ダイアログボックスが表示されます。 表示された [サービス認証] ダイアログボックスで、[OK] をクリックします。 次に、以下のパラメーターを設定し、作成 をクリックします。 サービスリンクロールがすでに作成されている場合、それらは再度作成されず、[サービス認証] ダイアログボックスは表示されません。
パラメーター | 説明 | 例 |
接続アドレス | MaxCompute のサービスエンドポイント。 エンドポイントのリージョン ID を MaxCompute プロジェクトのリージョン ID に置き換えます。 詳細については、「エンドポイント」をご参照ください。 - VPC エンドポイント:低レイテンシーで推奨されます。 ApsaraMQ for Kafka インスタンスと MaxCompute が同じリージョンにある場合に使用します。 - パブリックエンドポイント:高レイテンシーで推奨されません。 ApsaraMQ for Kafka インスタンスと MaxCompute が異なるリージョンにある場合に使用します。 パブリックエンドポイントを使用するには、コネクタのパブリックネットワークアクセスを有効にする必要があります。 詳細については、「コネクタのパブリックネットワークアクセスを有効にする」をご参照ください。 | http://service.cn-hangzhou.maxcompute.aliyun-inc.com/api |
ワークスペース | MaxCompute のワークスペース。宛先テーブルを含む MaxCompute プロジェクトに対応します。 | connector_test |
テーブル | MaxCompute のテーブル。 | test_kafka |
テーブルリージョン | MaxCompute テーブルが存在するリージョン。 | 中国 (杭州) |
サービスアカウント | MaxCompute の Alibaba Cloud アカウント ID。 | 188*** |
権限が付与されたロール名 | ApsaraMQ for Kafka の RAM ロールの名前。 詳細については、「RAM ロールの作成」をご参照ください。 | AliyunKafkaMaxComputeUser1 |
モード | メッセージがコネクタに同期されるモード。 デフォルト値:DEFAULT。 有効な値: - KEY:メッセージのキーのみが保持され、MaxCompute テーブルのキー列に書き込まれます。 - VALUE:メッセージの値のみが保持され、MaxCompute テーブルの値列に書き込まれます。 - DEFAULT:メッセージのキーと値の両方が保持され、MaxCompute テーブルのキー列と値列に書き込まれます。 > 重要: DEFAULT モードでは、CSV 形式はサポートされていません。 TEXT 形式と BINARY 形式のみがサポートされています。 | DEFAULT |
フォーマット | メッセージがコネクタに同期されるフォーマット。 デフォルト値:TEXT。 有効な値: - TEXT:メッセージは文字列です。 - BINARY:メッセージはバイト配列です。 - CSV:メッセージはカンマ (,) で区切られた文字列です。 > 重要: CSV 形式では、DEFAULT モードはサポートされていません。 KEY モードと VALUE モードのみがサポートされています: - KEY モード:メッセージのキーのみが保持されます。 キー文字列はカンマ (,) で区切られ、区切られた文字列はインデックスの順にテーブルに書き込まれます。 - VALUE モード:メッセージの値のみが保持されます。 値文字列はカンマ (,) で区切られ、区切られた文字列はインデックスの順にテーブルに書き込まれます。 | TEXT |
パーティション | パーティションの粒度。 デフォルト値:HOUR。 有効な値: - DAY:データは毎日新しいパーティションに書き込まれます。 - HOUR:データは毎時新しいパーティションに書き込まれます。 - MINUTE:データは毎分新しいパーティションに書き込まれます。 | HOUR |
タイムゾーン | コネクタのソーストピックにメッセージを送信する ApsaraMQ for Kafka プロデューサークライアントのタイムゾーン。 デフォルト値:GMT+08:00。 | GMT+08:00 |
コネクタが作成されると、Connector タスクリスト ページで表示できます。
コネクタが作成された後、Connector タスクリスト ページに移動し、作成したコネクタを見つけて、操作 列の デプロイ をクリックします。
テストメッセージの送信
MaxCompute sink コネクタをデプロイした後、ApsaraMQ for Kafka のソーストピックにメッセージを送信して、データが MaxCompute に同期できるかどうかをテストできます。
Connector タスクリスト ページで、対象のコネクタの 操作 列にある テスト をクリックします。
Connector タスクリスト ページで、目的のコネクタを探し、操作 列の テスト をクリックします。
Connector タスクリスト ページで、目的のコネクタを見つけ、操作 列の テスト をクリックします。
Connector タスクリスト ページで、目的のコネクタを見つけ、操作 列の テスト をクリックします。
Connector タスクリスト ページで、対象のコネクタを見つけ、操作 列の テスト をクリックします。
1. メッセージの送信 パネルで、テストメッセージを送信します。
送信方法 を コンソール に設定します。
メッセージキー テキストボックスに、メッセージのキー (たとえば demo) を入力します。
メッセージの内容 フィールドに、テストメッセージの内容 (例: {"key": "test"}) を入力します。
メッセージを特定のパーティションに送信するかどうかを指定するには、指定されたパーティションに送信 を設定します。
はい をクリックし、パーティション ID フィールドにパーティション ID (たとえば 0) を入力します。 パーティション ID のクエリ方法の詳細については、「パーティションステータスの表示」をご参照ください。
いいえ をクリックして、パーティションを指定せずにメッセージを送信します。
送信方法 を Docker に設定し、Docker コンテナーを実行してサンプルメッセージを生成する セクションで Docker コマンドを実行してメッセージを送信します。
送信方法 を SDK に設定します。 業務要件に応じて、使用するプログラミング言語またはフレームワークのソフトウェア開発キット (SDK) と接続方法を選択し、その SDK を使用してメッセージを送信します。
テーブルデータの表示
ApsaraMQ for Kafka のソーストピックにメッセージを送信した後、MaxCompute クライアントでテーブルデータを表示して、メッセージが受信されたかどうかを確認します。
以下の手順では、このトピックで `test_kafka` に書き込まれたデータを表示する方法を示します。 コネクタは、[パーティション] パラメーターに設定した粒度で、メッセージが書き込まれた時間に対応するパーティションに各メッセージを書き込みます。
MaxCompute クライアントにログインします。
次のコマンドを実行して、テーブルのデータパーティションを表示します。
show partitions test_kafka;次の結果が返されます。
pt=11-17-2020 15 OK次のコマンドを実行して、前の手順で返されたパーティションのデータを表示します。
select * from test_kafka where pt ="11-17-2020 15";次の結果が返されます。
+----------------------+------------+------------+-----+-------+---------------+
| topic | partition | offset | key | value | pt |
+----------------------+------------+------------+-----+-------+---------------+
| maxcompute-test-input| 0 | 0 | 1 | 1 | 11-17-2020 15 |
+----------------------+------------+------------+-----+-------+---------------+クエリで行が返されなかった場合は、show partitions が返したパーティションをクエリしたことを確認し、ソース Topic をサブスクライブする グループ の消費の詳細を確認し、コネクタのエラーデータ Topic とデッドレターキュー Topic に書き込みに失敗したメッセージがないかを確認してください。