Kafka データソースは、Kafka との間で双方向のデータ読み書きチャネルを提供します。このトピックでは、DataWorks が Kafka に対して提供するデータ同期機能について説明します。
サポートされるバージョン
DataWorks は、Alibaba Cloud Kafka および 0.10.2 から 3.6.x までのセルフマネージド Kafka バージョンをサポートしています。
0.10.2 より前のバージョンの Kafka は、パーティションオフセットの取得をサポートしておらず、データ構造がタイムスタンプをサポートしていない可能性があるため、データ同期はサポートされていません。
リアルタイム読み取り
-
Serverless リソースグループ (サブスクリプション) を使用する場合、リソース不足によるタスクの失敗を防ぐために、必要な仕様を事前に見積もる必要があります。
トピックごとに 1 CU を見積もります。トラフィックに基づいてリソースを見積もる必要もあります:
-
非圧縮の Kafka データの場合、10 MB/s のトラフィックごとに 1 CU を見積もります。
-
圧縮された Kafka データの場合、10 MB/s のトラフィックごとに 2 CU を見積もります。
-
JSON 解析が必要な圧縮された Kafka データの場合、10 MB/s のトラフィックごとに 3 CU を見積もります。
-
-
Serverless リソースグループ (サブスクリプション) または旧バージョンのデータ統合専用リソースグループを使用する場合:
-
ご利用のワークロードのフェールオーバーに対する許容度が高い場合、クラスターのスロット使用率は 80% を超えないようにしてください。
-
ご利用のワークロードのフェールオーバーに対する許容度が低い場合、クラスターのスロット使用率は 70% を超えないようにしてください。
-
実際のリソース使用量は、データの内容やフォーマットなどの要因によって異なります。最初のリソース評価の後、実際のランタイム使用量に基づいてリソースを調整してください。
制限事項
Kafka データソースは、Serverless リソースグループ (推奨) および 旧バージョンのデータ統合専用リソースグループ をサポートしています。
単一テーブルからのオフライン読み取り
parameter.groupId と parameter.kafkaConfig.group.id の両方が設定されている場合、parameter.groupId が kafkaConfig パラメーター内の group.id よりも優先されます。
単一テーブルへのリアルタイム書き込み
書き込み操作では、データの重複排除はサポートされていません。オフセットのリセットまたはフェールオーバーの後にタスクが再起動されると、重複したデータが書き込まれる可能性があります。
データベース全体のリアルタイム書き込み
-
リアルタイムデータ同期タスクは、Serverless リソースグループ (推奨) および 旧バージョンのデータ統合専用リソースグループ をサポートしています。
-
ソーステーブルにプライマリキーがある場合、プライマリキーの値が Kafka レコードのキーとして使用されます。これにより、同じプライマリキーへの変更が同じ Kafka パーティションに順序通りに書き込まれることが保証されます。
-
ソーステーブルにプライマリキーがない場合は、2 つのオプションがあります。プライマリキーなしでテーブルを同期するオプションを選択した場合、Kafka レコードのキーは空になります。テーブルの変更が Kafka に順序通りに書き込まれるようにするには、送信先の Kafka トピックに単一のパーティションのみを設定する必要があります。カスタムプライマリキーを選択した場合、1 つ以上の非プライマリキーフィールドの組み合わせが Kafka レコードのキーとして使用されます。
-
Kafka クラスターが例外を返した場合でも、同じプライマリキーの変更が同じ Kafka パーティションに順序通りに書き込まれるようにするには、拡張パラメーターフォームに次の設定を追加します。
{"max.in.flight.requests.per.connection":1,"buffer.memory": 100554432}重要この設定は、同期性能を大幅に低下させます。性能と、厳密な順序付けおよび信頼性の必要性とのバランスを取る必要があります。
-
リアルタイム同期で Kafka に書き込まれるメッセージの全体的なフォーマット、ハートビートメッセージのフォーマット、およびソースデータの変更に対応するメッセージのフォーマットの詳細については、「付録:メッセージフォーマット」をご参照ください。
サポートされるフィールドタイプ
Kafka は非構造化データストレージを提供します。Kafka レコードには通常、key、value、offset、timestamp、headers、partition などのフィールドが含まれます。DataWorks が Kafka からデータを読み書きする際、次のようにデータを処理します。
データの読み取り
DataWorks が Kafka からデータを読み取る際、データを JSON フォーマットで解析できます。次の表に、各データモジュールの処理方法を示します。
|
Kafka レコードデータモジュール |
処理されるデータ型 |
|
key |
データ同期タスクの keyType 設定項目に依存します。keyType パラメーターの詳細については、付録の完全なパラメーター説明をご参照ください。 |
|
value |
データ同期タスクの valueType 設定項目に依存します。valueType パラメーターの詳細については、付録の完全なパラメーター説明をご参照ください。 |
|
offset |
Long |
|
timestamp |
Long |
|
headers |
String |
|
partition |
Long |
データの書き込み
DataWorks は、JSON またはテキスト形式で Kafka にデータを書き込みます。処理ポリシーは、以下の表に示すように同期タスクの種類によって異なります。
-
データがテキスト形式で書き込まれる場合、フィールド名は含まれません。フィールド値は区切り文字で区切られます。
-
リアルタイム同期タスクが Kafka にデータを書き込む場合、組み込みの JSON 形式を使用します。データには、データベースの変更メッセージ、業務日時、データ定義言語 (DDL) 情報などが含まれます。データ形式の詳細については、「付録:メッセージフォーマット」をご参照ください。
|
同期タスクの種類 |
Kafka に書き込まれる value のフォーマット |
ソースフィールドの型 |
書き込み操作の処理方法 |
|
オフライン同期 DataStudio のオフライン同期ノード |
json |
String |
UTF-8 エンコードされた文字列 |
|
Boolean |
UTF-8 エンコードされた文字列 "true" または "false" に変換 |
||
|
Time/Date |
yyyy-MM-dd HH:mm:ss 形式の UTF-8 エンコードされた文字列 |
||
|
Numeric |
UTF-8 エンコードされた数値文字列 |
||
|
Byte stream |
バイトストリームは UTF-8 エンコードされた文字列として扱われ、文字列に変換されます。 |
||
|
text |
String |
UTF-8 エンコードされた文字列 |
|
|
Boolean |
UTF-8 エンコードされた文字列 "true" または "false" に変換 |
||
|
Time/Date |
yyyy-MM-dd HH:mm:ss 形式の UTF-8 エンコードされた文字列 |
||
|
Numeric |
UTF-8 エンコードされた数値文字列 |
||
|
Byte stream |
バイトストリームは UTF-8 エンコードされた文字列として扱われ、文字列に変換されます。 |
||
|
リアルタイム同期:Kafka へのリアルタイム ETL DataStudio のリアルタイム同期ノード |
json |
String |
UTF-8 エンコードされた文字列 |
|
Boolean |
JSON ブール型 |
||
|
Time/Date |
|
||
|
Numeric |
JSON 数値型 |
||
|
Byte stream |
バイトストリームは Base64 エンコードされた後、UTF-8 エンコードされた文字列に変換されます。 |
||
|
text |
String |
UTF-8 エンコードされた文字列 |
|
|
Boolean |
UTF-8 エンコードされた文字列 "true" または "false" に変換 |
||
|
Time/Date |
yyyy-MM-dd HH:mm:ss 形式の UTF-8 エンコードされた文字列 |
||
|
Numeric |
UTF-8 エンコードされた数値文字列 |
||
|
Byte stream |
バイトストリームは Base64 エンコードされた後、UTF-8 エンコードされた文字列に変換されます。 |
||
|
リアルタイム同期:データベース全体を Kafka にリアルタイム同期 増分データのみのリアルタイム同期 |
組み込み JSON 形式 |
String |
UTF-8 エンコードされた文字列 |
|
Boolean |
JSON ブール型 |
||
|
Time/Date |
13 桁のミリ秒タイムスタンプ |
||
|
Numeric |
JSON 数値 |
||
|
Byte stream |
バイトストリームは Base64 エンコードされた後、UTF-8 エンコードされた文字列に変換されます。 |
||
|
同期ソリューション:ワンクリックで Kafka にリアルタイム同期 完全オフライン同期 + 増分リアルタイム同期 |
組み込み JSON 形式 |
String |
UTF-8 エンコードされた文字列 |
|
Boolean |
JSON ブール型 |
||
|
Time/Date |
13 桁のミリ秒タイムスタンプ |
||
|
Numeric |
JSON 数値 |
||
|
Byte stream |
バイトストリームは Base64 エンコードされた後、UTF-8 エンコードされた文字列に変換されます。 |
データソースの追加
DataWorks で同期タスクを開発する前に、「データソース設定」の指示に従って、必要なデータソースを DataWorks に追加する必要があります。データソースを追加する際に、DataWorks コンソールでパラメーターの説明を表示して、パラメーターの意味を理解できます。
データ同期タスクの開発
同期タスクの設定のエントリポイントと手順については、以下の設定ガイドをご参照ください。
単一テーブルのオフライン同期タスクの設定
-
手順の詳細については、「コードレス UI でオフライン同期タスクを設定する」および「コードエディタでオフライン同期タスクを設定する」をご参照ください。
-
コードエディタのパラメーターの完全なリストとスクリプトのデモについては、「付録:スクリプトのデモとパラメーターの説明」をご参照ください。
単一テーブルまたはデータベース全体のリアルタイム同期タスクの設定
手順については、「単一テーブルのリアルタイム同期タスクを設定する」および「データベース全体のリアルタイム同期タスクを設定する」をご参照ください。
認証設定
SSL
Kafka データソースを設定する際に、[特別な認証方法] を [SSL] または [SASL_SSL] に設定すると、Kafka クラスターの SSL 認証が有効になります。クライアントのトラストストア証明書ファイルをアップロードし、トラストストアのパスフレーズを入力する必要があります。
-
Kafka クラスターが Alibaba Cloud Kafka インスタンスの場合、「SSL 証明書アルゴリズムのアップグレード手順」を参照して、正しいトラストストア証明書ファイルをダウンロードしてください。トラストストアのパスフレーズは KafkaOnsClient です。
-
Kafka クラスターが EMR インスタンスの場合、「Kafka 接続に SSL 暗号化を使用する」を参照して、正しいトラストストア証明書ファイルをダウンロードし、トラストストアのパスフレーズを取得してください。
-
セルフマネージドクラスターの場合、正しいトラストストア証明書をアップロードし、正しいトラストストアのパスフレーズを入力する必要があります。
キーストア証明書ファイル、キーストアのパスフレーズ、および SSL パスフレーズは、Kafka クラスターで双方向 SSL 認証が有効になっている場合にのみ必要です。Kafka クラスターサーバーはこれを使用してクライアントの ID を認証します。双方向 SSL 認証は、Kafka クラスターの server.properties ファイルで ssl.client.auth=required が設定されている場合に有効になります。詳細については、「Kafka 接続に SSL 暗号化を使用する」をご参照ください。
GSSAPI
Kafka データソースを設定する際に [Sasl メカニズム] を GSSAPI に設定した場合、3 つの認証ファイルをアップロードする必要があります:JAAS 設定ファイル、Kerberos 設定ファイル、および Keytab ファイル。また、専用リソースグループの DNS/HOST 設定も構成する必要があります。以下のセクションでは、これらのファイルと必要な DNS および HOST 設定について説明します。
Serverless リソースグループ の場合、内部 DNS 名前解決を使用してホストアドレス情報を設定する必要があります。詳細については、「内部 DNS 名前解決 (PrivateZone)」をご参照ください。
-
JAAS 設定ファイル
JAAS ファイルは KafkaClient で始まり、その後にすべての設定項目を中括弧 {} で囲む必要があります:
-
中括弧内の最初の行は、使用するログオンコンポーネントクラスを定義します。異なる SASL 認証メカニズムでは、ログオンコンポーネントクラスは固定されています。後続の各設定項目は key=value 形式で記述されます。
-
最後の設定項目を除き、すべての設定項目はセミコロンで終わってはいけません。
-
最後の設定項目はセミコロンで終わり、閉じ括弧
}の後にもう 1 つセミコロンを付ける必要があります。
フォーマット要件が満たされていない場合、JAAS 設定ファイルは解析できません。次のコードは、典型的な JAAS 設定ファイルのフォーマットを示しています。xxx プレースホルダーを実際の情報に置き換えてください。
KafkaClient { com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true keyTab="xxx" storeKey=true serviceName="kafka-server" principal="kafka-client@EXAMPLE.COM"; };設定項目
説明
ログオンモジュール
com.sun.security.auth.module.Krb5LoginModule に設定する必要があります。
useKeyTab
true に設定する必要があります。
keyTab
任意のパスを指定できます。同期タスクの実行時に、システムはデータソース設定時にアップロードした keytab ファイルをローカルパスに自動的にダウンロードします。その後、システムはこのローカルパスを keytab 設定項目に使用します。
storeKey
クライアントがキーを保存するかどうかを指定します。true または false に設定できます。データ同期には影響しません。
serviceName
Kafka サーバーの server.properties 設定ファイル内の sasl.kerberos.service.name 設定項目に対応します。必要に応じてこの項目を設定してください。
principal
Kafka クライアントが使用する Kerberos プリンシパルです。必要に応じて設定し、アップロードした keytab ファイルにこのプリンシパルのキーが含まれていることを確認してください。
-
-
Kerberos 設定ファイル
Kerberos 設定ファイルには、[libdefaults] と [realms] の 2 つのモジュールが含まれている必要があります。
-
[libdefaults] モジュールは Kerberos 認証パラメーターを指定します。モジュール内の各設定項目は key=value 形式で記述されます。
-
[realms] モジュールはキー配布センター (KDC) のアドレスを指定します。複数のレルムサブモジュールを含むことができます。各レルムサブモジュールは、レルム名の後に等号 (=) が続きます。
その後に、中括弧で囲まれた設定項目のセットが続きます。各設定項目も key=value 形式で記述されます。次のコードは、典型的な Kerberos 設定ファイルのフォーマットを示しています。xxx プレースホルダーを実際の情報に置き換えてください。
[libdefaults] default_realm = xxx [realms] xxx = { kdc = xxx }設定項目
説明
[libdefaults].default_realm
Kafka クラスターノードにアクセスする際に使用されるデフォルトのレルムです。これは通常、JAAS 設定ファイルで指定されたクライアントプリンシパルのレルムと同じです。
その他の [libdefaults] パラメーター
[libdefaults] モジュールは、ticket_lifetime などの他の Kerberos 認証パラメーターを指定できます。必要に応じて設定してください。
[realms].realm name
JAAS 設定ファイルで指定されたクライアントプリンシパルのレルムおよび [libdefaults].default_realm と同じである必要があります。JAAS 設定ファイルのクライアントプリンシパルのレルムが [libdefaults].default_realm と異なる場合は、2 つのレルムサブモジュールを含める必要があります。これらのサブモジュールは、それぞれ JAAS 設定ファイルのクライアントプリンシパルのレルムと [libdefaults].default_realm に対応している必要があります。
[realms].realm name.kdc
KDC のアドレスとポートを ip:port 形式で指定します。例:kdc=10.0.0.1:88。ポートを省略した場合、システムはデフォルトポート 88 を使用します。例:kdc=10.0.0.1。
-
-
Keytab ファイル
keytab ファイルには、JAAS 設定ファイルで指定されたプリンシパルのキーが含まれている必要があり、KDC によって検証可能である必要があります。たとえば、現在の作業ディレクトリに client.keytab という名前のファイルがある場合、次のコマンドを実行して、keytab ファイルに指定されたプリンシパルのキーが含まれているかどうかを確認できます。
klist -ket ./client.keytab Keytab name: FILE:client.keytab KVNO Timestamp Principal ---- ------------------- ------------------------------------------------------ 7 2018-07-30T10:19:16 te**@**.com (des-cbc-md5) -
専用リソースグループの DNS および HOST 設定
Kafka クラスターが Kerberos 認証を使用する場合、KDC は各ノードのホスト名を使用して各ノードのプリンシパルを登録します。クライアントが Kafka クラスターノードに接続すると、ローカルの DNS および HOST 設定を使用してノードのプリンシパルを導き出し、KDC にノードのアクセス認証情報を要求します。Kerberos 認証が有効になっている Kafka クラスターにアクセスするために専用リソースグループを使用する場合、クラスターノードのアクセス認証情報を KDC から取得できるように、DNS および HOST 設定を正しく構成する必要があります:
-
DNS 設定
専用リソースグループがアタッチされている VPC 内の Kafka クラスターノードの名前解決に PrivateZone インスタンスを使用する場合、IP アドレス 100.100.2.136 および 100.100.2.138 のカスタムルートを VPC アタッチメントに追加できます。これにより、Kafka クラスターノードの PrivateZone 名前解決設定が専用リソースグループに適用されることが保証されます。DataWorks コンソールの左側のナビゲーションウィンドウで、[リソースグループリスト] をクリックします。ご利用の専用リソースグループの [アクション] 列で、[ネットワーク設定] をクリックします。[VPC アタッチメント] タブで、[アクション] 列の [カスタムルート] をクリックします。表示されるダイアログボックスで、[ルートの追加] をクリックし、[宛先タイプ] を [IDC] に、[接続方法] を [直接 IP] に設定し、直接 IP アドレスを入力してから、[ルートの生成] をクリックします。
-
HOST 設定
専用リソースグループがアタッチされている VPC 内の Kafka クラスターノードの名前解決に PrivateZone インスタンスを使用しない場合は、各 Kafka クラスターノードの IP アドレスとドメイン名のマッピングをホスト設定に追加する必要があります。DataWorks コンソールの左側のナビゲーションウィンドウで、[リソースグループリスト] をクリックします。ご利用のリソースグループを見つけ、[アクション] 列の [ネットワーク設定] をクリックし、[ホスト設定] タブを選択します。[追加] をクリックしてホストドメインを追加します。ホスト設定は DNS 設定よりも優先されます。
-
PLAIN
Kafka データソースを設定する際に [Sasl メカニズム] を PLAIN に設定した場合、JAAS ファイルは KafkaClient で始まり、その後にすべての設定項目を中括弧 {} で囲む必要があります。
-
中括弧内の最初の行は、使用するログオンコンポーネントクラスを定義します。異なる SASL 認証メカニズムでは、ログオンコンポーネントクラスは固定されています。後続の各設定項目は key=value 形式で記述されます。
-
最後の設定項目を除き、すべての設定項目はセミコロンで終わってはいけません。
-
最後の設定項目はセミコロンで終わる必要があります。閉じ括弧 "}" の後にもセミコロンを追加する必要があります。
フォーマット要件が満たされていない場合、JAAS 設定ファイルは解析できません。次のコードは、典型的な JAAS 設定ファイルのフォーマットを示しています。xxx プレースホルダーを実際の情報に置き換えてください。
KafkaClient {
org.apache.kafka.common.security.plain.PlainLoginModule required
username="xxx"
password="xxx";
};
|
設定項目 |
説明 |
|
ログオンモジュール |
org.apache.kafka.common.security.plain.PlainLoginModul に設定する必要があります |
|
username |
ユーザー名です。必要に応じてこの項目を設定してください。 |
|
password |
パスワードです。必要に応じてこの項目を設定してください。 |
よくある質問
付録:スクリプトのデモとパラメーターの説明
コードエディタを使用したバッチ同期タスクの設定
コードエディタを使用してバッチ同期タスクを設定する場合、統一されたスクリプトフォーマット要件に基づいて、スクリプト内の関連パラメーターを設定する必要があります。詳細については、「スクリプトモードの設定」をご参照ください。以下の情報は、コードエディタを使用してバッチ同期タスクを設定する際に、データソースに対して設定する必要があるパラメーターについて説明しています。
Reader スクリプトのデモ
次の JSON 設定は、Kafka からデータを読み取ります。
{
"type": "job",
"steps": [
{
"stepType": "kafka",
"parameter": {
"server": "host:9093",
"column": [
"__key__",
"__value__",
"__partition__",
"__offset__",
"__timestamp__",
"'123'",
"event_id",
"tag.desc"
],
"kafkaConfig": {
"group.id": "demo_test"
},
"topic": "topicName",
"keyType": "ByteArray",
"valueType": "ByteArray",
"beginDateTime": "20190416000000",
"endDateTime": "20190416000006",
"skipExceedRecord": "true"
},
"name": "Reader",
"category": "reader"
},
{
"stepType": "stream",
"parameter": {
"print": false,
"fieldDelimiter": ","
},
"name": "Writer",
"category": "writer"
}
],
"version": "2.0",
"order": {
"hops": [
{
"from": "Reader",
"to": "Writer"
}
]
},
"setting": {
"errorLimit": {
"record": "0"
},
"speed": {
"throttle": true,//throttle が false の場合、mbps パラメーターは有効にならず、レートは制限されません。throttle が true の場合、レートは制限されます。
"concurrent": 1,//並行スレッドの数。
"mbps":"12"//最大転送レート。1 mbps は 1 MB/s に相当します。
}
}
}
Reader スクリプトのパラメーター
|
パラメーター |
説明 |
必須 |
|
datasource |
データソースの名前です。コードエディタはデータソースの追加をサポートしています。このパラメーターの値は、追加されたデータソースの名前と同じである必要があります。 |
はい |
|
server |
Kafka ブローカーサーバーのアドレスを ip:port 形式で指定します。 server は 1 つしか設定できませんが、DataWorks が Kafka クラスター内のすべてのブローカーの IP アドレスに接続できることを確認する必要があります。 |
はい |
|
topic |
Kafka トピックです。トピックは、Kafka が処理するメッセージフィードの集約です。 |
はい |
|
column |
読み取る Kafka データです。定数列、データ列、属性列がサポートされています。
|
はい |
|
keyType |
Kafka キーの型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。 |
いいえ |
|
valueType |
Kafka 値の型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。 |
いいえ |
|
beginDateTime |
データ消費の開始時刻です。このパラメーターは時間範囲の左境界を指定し、その時刻を含みます。yyyymmddhhmmss 形式の時刻文字列です。このパラメーターは スケジューリングパラメーター と共に使用できます。詳細については、「サポートされるスケジューリングパラメーターの形式」をご参照ください。 説明
この機能は Kafka 0.10.2 以降でサポートされています。 |
このパラメーターまたは beginOffset のいずれかを指定する必要があります。 説明
beginDateTime と endDateTime は一緒に使用されます。 |
|
endDateTime |
データ消費の終了時刻です。このパラメーターは時間範囲の右境界を指定し、その時刻を含みません。yyyymmddhhmmss 形式の時刻文字列です。このパラメーターは スケジューリングパラメーター と共に使用できます。詳細については、「サポートされるスケジューリングパラメーターの形式」をご参照ください。 説明
この機能は Kafka 0.10.2 以降でサポートされています。 |
このパラメーターまたは endOffset のいずれかを指定する必要があります。 説明
endDateTime と beginDateTime は一緒に使用されます。 |
|
beginOffset |
データ消費の開始オフセットです。次の形式で設定できます:
|
このパラメーターまたは beginDateTime のいずれかを指定する必要があります。 |
|
endOffset |
データ消費の終了オフセットです。これは、データ消費タスクがいつ終了するかを制御するために使用されます。 |
このパラメーターまたは endDateTime のいずれかを指定する必要があります。 |
|
skipExceedRecord |
Kafka は
|
いいえ。デフォルト値は false です。 |
|
partition |
Kafka トピックには複数のパーティション (partition) があります。デフォルトでは、データ同期タスクはトピック内のすべてのパーティションをカバーするオフセット範囲からデータを読み取ります。partition を指定して、単一のパーティションのオフセット範囲からのみデータを読み取ることもできます。 |
いいえ。デフォルト値はありません。 |
|
kafkaConfig |
データ消費のために KafkaConsumer クライアントを作成する際に、bootstrap.servers、auto.commit.interval.ms、session.timeout.ms などの拡張パラメーターを指定できます。kafkaConfig を使用して、KafkaConsumer の消費動作を制御できます。 |
いいえ |
|
encoding |
keyType または valueType が STRING に設定されている場合、このパラメーターで指定されたエンコーディングが文字列の解析に使用されます。 |
いいえ。デフォルト値は UTF-8 です。 |
|
waitTIme |
コンシューマーオブジェクトが 1 回の試行で Kafka からデータをプルするのを待つ最大時間 (秒)。 |
いいえ。デフォルト値は 60 です。 |
|
stopWhenPollEmpty |
有効な値は true と false です。このパラメーターが true に設定され、コンシューマーが Kafka から空のデータをプルした場合 (通常はトピック内のすべてのデータが読み取られたため、またはネットワークや Kafka クラスターの可用性の問題のため)、タスクはすぐに停止します。それ以外の場合は、データが再び読み取られるまで再試行します。 |
いいえ。デフォルト値は true です。 |
|
stopWhenReachEndOffset |
このパラメーターは、stopWhenPollEmpty が true の場合にのみ有効です。有効な値は true と false です。
|
いいえ。デフォルト値は false です。 説明
この設定は下位互換性を提供します。0.10.2 より前の Kafka バージョンは、すべてのパーティションの最大オフセットのチェックをサポートしていません。 |
次の表に、kafkaConfig パラメーターを示します。
|
パラメーター |
説明 |
|
fetch.min.bytes |
コンシューマーが 1 回のリクエストでブローカーからフェッチするデータの最小量 (バイト)。ブローカーは、この量のデータが利用可能になるまで待機してからコンシューマーに応答します。 |
|
fetch.max.wait.ms |
ブローカーがフェッチリクエストに応答する前にデータが利用可能になるのを待つ最大時間 (ミリ秒)。デフォルト値は 500 です。ブローカーは、fetch.min.bytes または fetch.max.wait.ms のいずれかの条件が満たされたときに応答します。 |
|
max.partition.fetch.bytes |
ブローカーが各 パーティション からコンシューマーに返すことができる最大バイト数を指定します。デフォルト値は 1 MB です。 |
|
session.timeout.ms |
コンシューマーがサービスを受け取らなくなる前にサーバーから切断できる時間を指定します。デフォルト値は 30 秒です。 |
|
auto.offset.reset |
オフセットなし、または無効なオフセット (コンシューマーが長時間非アクティブで、オフセットを持つレコードが期限切れで削除されたため) で読み取る際にコンシューマーが実行するアクション。デフォルト値は none で、オフセットは自動的にリセットされません。earliest に変更でき、これはコンシューマーが パーティション レコードを最小オフセットから読み取ることを意味します。 |
|
max.poll.records |
poll メソッドの 1 回の呼び出しで返すことができるメッセージの数。 |
|
key.deserializer |
メッセージキーの逆シリアル化メソッド。例:org.apache.kafka.common.serialization.StringDeserializer。 |
|
value.deserializer |
データ値の逆シリアル化メソッド。例:org.apache.kafka.common.serialization.StringDeserializer。 |
|
ssl.truststore.location |
SSL ルート証明書のパス。 |
|
ssl.truststore.password |
ルート証明書ストアのパスワード。Alibaba Cloud Kafka を使用している場合は、これを KafkaOnsClient に設定します。 |
|
security.protocol |
アクセスプロトコル。現在、SASL_SSL プロトコルのみがサポートされています。 |
|
sasl.mechanism |
SASL 認証方式。Alibaba Cloud Kafka を使用している場合は、PLAIN を使用します。 |
|
java.security.auth.login.config |
SASL 認証ファイルのパス。 |
Writer スクリプトのデモ
次のコードは、Kafka にデータを書き込むための JSON 設定を示しています。
{
"type":"job",
"version":"2.0",//バージョン番号。
"steps":[
{
"stepType":"stream",
"parameter":{},
"name":"Reader",
"category":"reader"
},
{
"stepType":"Kafka",//プラグイン名。
"parameter":{
"server": "ip:9092", //Kafka のサーバーアドレス。
"keyIndex": 0, //キーとして使用する列。キャメルケースの命名規則に従い、k は小文字にする必要があります。
"valueIndex": 1, //値として使用する列。現在、ソースデータから 1 つの列を選択するか、このパラメーターを空にすることができます。空の場合、すべてのソースデータが使用されます。
//たとえば、ODPS テーブルの 2 番目、3 番目、4 番目の列を kafkaValue として使用するには、新しい ODPS テーブルを作成し、元の ODPS テーブルのデータを新しいテーブルにクリーンアップして統合し、その後、新しいテーブルを同期に使用します。
"keyType": "Integer", //Kafka キーの型。
"valueType": "Short", //Kafka 値の型。
"topic": "t08", //Kafka トピック。
"batchSize": 1024 //一度に Kafka に書き込むデータ量 (バイト)。
},
"name":"Writer",
"category":"writer"
}
],
"setting":{
"errorLimit":{
"record":"0"//エラーレコードの数。
},
"speed":{
"throttle":true,//throttle が false の場合、mbps パラメーターは有効にならず、レートは制限されません。throttle が true の場合、レートは制限されます。
"concurrent":1, //同時実行ジョブの数。
"mbps":"12"//最大転送レート。1 mbps は 1 MB/s に相当します。
}
},
"order":{
"hops":[
{
"from":"Reader",
"to":"Writer"
}
]
}
}
Writer スクリプトのパラメーター
|
パラメーター |
説明 |
必須 |
|
datasource |
データソースの名前です。コードエディタはデータソースの追加をサポートしています。このパラメーターの値は、追加されたデータソースの名前と同じである必要があります。 |
はい |
|
server |
Kafka サーバーのアドレスを ip:port 形式で指定します。 |
はい |
|
topic |
Kafka トピックです。Kafka が処理するさまざまなメッセージフィードのカテゴリです。 Kafka クラスターに公開される各メッセージにはカテゴリがあり、これをトピックと呼びます。トピックはメッセージのグループのコレクションです。 |
はい |
|
valueIndex |
Kafka ライターで値として使用される列です。これを指定しない場合、デフォルトですべての列が連結されて値を形成します。区切り文字は fieldDelimiter で指定されます。 |
いいえ |
|
writeMode |
valueIndex が設定されていない場合、このパラメーターはソースレコードのすべての列を連結して Kafka レコードの値を形成するフォーマットを決定します。有効な値は text と JSON です。デフォルト値は text です。
たとえば、ソースレコードに a、b、c の値を持つ 3 つの列があり、writeMode が text に、fieldDelimiter が # に設定されている場合、書き込まれる Kafka レコードの値は文字列 a#b#c です。writeMode が JSON に設定され、column が [{"name":"col1"},{"name":"col2"},{"name":"col3"}] に設定されている場合、書き込まれる Kafka レコードの値は文字列 {"col1":"a","col2":"b","col3":"c"} です。 valueIndex が設定されている場合、このパラメーターは無効です。 |
いいえ |
|
column |
データを書き込む送信先テーブルのフィールドをカンマで区切って指定します。例: valueIndex が設定されておらず、writeMode が JSON に設定されている場合、このパラメーターはソースレコードの列値の JSON 構造内のフィールド名を定義します。例:
valueIndex が設定されているか、writeMode が text に設定されている場合、このパラメーターは無効です。 |
valueIndex が設定されておらず、writeMode が JSON に設定されている場合に必須です。 |
|
partition |
データが書き込まれる Kafka トピック内のパーティションの番号を指定します。これは 0 以上の整数でなければなりません。 |
いいえ |
|
keyIndex |
Kafka ライターでキーとして使用される列です。 keyIndex パラメーターの値は 0 以上の整数でなければなりません。そうでない場合、タスクは失敗します。 |
いいえ |
|
keyIndexes |
Kafka レコードのキーとして使用されるソースレコード内の列の序数の配列です。 列の序数は 0 から始まります。たとえば、[0,1,2] は、設定されたすべての列番号の値をカンマで連結して Kafka レコードのキーを形成します。これを指定しない場合、Kafka レコードのキーは null になり、データはトピックのパーティションにラウンドロビン方式で書き込まれます。このパラメーターまたは keyIndex のいずれか一方のみを指定できます。 |
いいえ |
|
fieldDelimiter |
writeMode が text に設定され、valueIndex が設定されていない場合、ソースレコードのすべての列がこのパラメーターで指定された列区切り文字を使用して連結され、Kafka レコードの値を形成します。区切り文字として単一の文字または複数の文字を設定できます。Unicode 文字は \u0001 の形式で設定できます。\t や \n などのエスケープ文字がサポートされています。デフォルト値は \t です。 writeMode が text に設定されていないか、valueIndex が設定されている場合、このパラメーターは無効です。 |
いいえ |
|
keyType |
Kafka キーの型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。 |
はい |
|
valueType |
Kafka 値の型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。 |
はい |
|
nullKeyFormat |
keyIndex または keyIndexes で指定されたソース列の値が null の場合、このパラメーターで指定された文字列に置き換えられます。これを設定しない場合、置き換えは行われません。 |
いいえ |
|
nullValueFormat |
ソース列の値が null の場合、Kafka レコードの値を組み立てる際に、このパラメーターで指定された文字列に置き換えられます。このパラメーターを指定しない場合、置き換えは行われません。 |
いいえ |
|
acks |
Kafka プロデューサーを初期化する際の acks 設定です。これは、書き込み成功の確認方法を決定します。デフォルトでは、acks パラメーターは all に設定されています。acks の有効な値は次のとおりです:
|
いいえ |
付録:Kafka への書き込みメッセージフォーマットの定義
リアルタイム同期タスクを設定して実行すると、ソースデータベースから読み取られたデータが JSON 形式で Kafka トピックに書き込まれます。まず、指定されたソーステーブル内のすべての既存データが対応する Kafka トピックに書き込まれます。その後、タスクはリアルタイム同期を開始し、増分データを継続的にトピックに書き込みます。ソーステーブルからの増分 DDL 変更情報も JSON 形式で Kafka トピックに書き込まれます。Kafka に書き込まれたメッセージのステータスと変更情報を取得できます。詳細については、「付録:メッセージフォーマット」をご参照ください。
オフライン同期タスクからの JSON データでは、payload.sequenceId、payload.timestamp.eventTime、および payload.timestamp.checkpointTime フィールドは -1 に設定されます。
付録:JSON フィールドタイプ
writeMode が JSON に設定されている場合、column パラメーター内の type フィールドを使用して JSON フィールドタイプを定義できます。書き込み操作中に、システムはソースレコードの列値を指定されたタイプに変換しようとします。タイプ変換に失敗すると、ダーティデータが発生します。
|
有効な値 |
説明 |
|
JSON_STRING |
ソースレコードの列値を文字列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が整数 |
|
JSON_NUMBER |
ソースレコードの列値を数値に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 |
|
JSON_BOOL |
ソースレコードの列値をブール値に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 |
|
JSON_ARRAY |
ソースレコードの列値を JSON 配列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 |
|
JSON_MAP |
ソースレコードの列値を JSON オブジェクトに変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 |
|
JSON_BASE64 |
ソース列のバイト配列を BASE64 エンコードされた文字列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が 16 進数で |
|
JSON_HEX |
ソース列のバイト配列を 16 進文字列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が 16 進数で |