すべてのプロダクト
Search
ドキュメントセンター

DataWorks:Kafka データソース

最終更新日:Aug 27, 2026

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.groupIdparameter.kafkaConfig.group.id の両方が設定されている場合、parameter.groupIdkafkaConfig パラメーター内の 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 レコードには通常、keyvalueoffsettimestampheaderspartition などのフィールドが含まれます。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

  • ミリ秒未満の精度の時刻値の場合:ミリ秒単位のタイムスタンプを表す 13 桁の JSON 整数に変換されます。

  • マイクロ秒またはナノ秒の精度の時刻値の場合:ミリ秒タイムスタンプ用の 13 桁の整数とナノ秒タイムスタンプ用の 6 桁の小数を含む JSON 浮動小数点数に変換されます。

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 コンソールでパラメーターの説明を表示して、パラメーターの意味を理解できます

データ同期タスクの開発

同期タスクの設定のエントリポイントと手順については、以下の設定ガイドをご参照ください。

単一テーブルのオフライン同期タスクの設定

単一テーブルまたはデータベース全体のリアルタイム同期タスクの設定

手順については、「単一テーブルのリアルタイム同期タスクを設定する」および「データベース全体のリアルタイム同期タスクを設定する」をご参照ください。

認証設定

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 データです。定数列、データ列、属性列がサポートされています。

  • 定数列:単一引用符で囲まれた列。例:["'abc'", "'123'"]

  • データ列

    • データが JSON 形式の場合、JSON オブジェクトのプロパティを取得できます。例:["event_id"]

    • データが JSON 形式の場合、JSON オブジェクトのネストされたサブプロパティを取得できます。例:["tag.desc"]

  • 属性列

    • __key__:メッセージのキー。

    • __value__:メッセージの完全な内容。

    • __partition__:現在のメッセージが存在するパーティション。

    • __headers__:現在のメッセージのヘッダー。

    • __offset__:現在のメッセージのオフセット。

    • __timestamp__:現在のメッセージのタイムスタンプ。

    次のコードは完全な例です。

    "column": [
        "__key__",
        "__value__",
        "__partition__",
        "__offset__",
        "__timestamp__",
        "'123'",
        "event_id",
        "tag.desc"
        ]

はい

keyType

Kafka キーの型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。

いいえ

valueType

Kafka 値の型です。有効な値:BYTEARRAY、DOUBLE、FLOAT、INTEGER、LONG、SHORT。

いいえ

beginDateTime

データ消費の開始時刻です。このパラメーターは時間範囲の左境界を指定し、その時刻を含みます。yyyymmddhhmmss 形式の時刻文字列です。このパラメーターは スケジューリングパラメーター と共に使用できます。詳細については、「サポートされるスケジューリングパラメーターの形式」をご参照ください。

説明

この機能は Kafka 0.10.2 以降でサポートされています。

このパラメーターまたは beginOffset のいずれかを指定する必要があります。

説明

beginDateTimeendDateTime は一緒に使用されます。

endDateTime

データ消費の終了時刻です。このパラメーターは時間範囲の右境界を指定し、その時刻を含みません。yyyymmddhhmmss 形式の時刻文字列です。このパラメーターは スケジューリングパラメーター と共に使用できます。詳細については、「サポートされるスケジューリングパラメーターの形式」をご参照ください。

説明

この機能は Kafka 0.10.2 以降でサポートされています。

このパラメーターまたは endOffset のいずれかを指定する必要があります。

説明

endDateTimebeginDateTime は一緒に使用されます。

beginOffset

データ消費の開始オフセットです。次の形式で設定できます:

  • 数値、例:15553274。これは消費の開始オフセットを示します。

  • seekToBeginning:最小オフセットからデータを消費することを示します。

  • seekToLastkafkaConfig パラメーターの group.id で指定されたグループ ID に保存されているオフセットからデータの読み取りを開始します。グループオフセットはクライアントによって定期的に Kafka サーバーに自動的にコミットされることに注意してください。そのため、タスクが失敗して再実行されると、データの重複や損失が発生する可能性があります。skipExceedRecord パラメーターが true に設定されている場合、タスクは読み取られた最後のいくつかのレコードを破棄する可能性があります。この破棄されたデータのグループオフセットはすでにサーバーにコミットされているため、このデータは次のタスク実行では読み取れません。

  • seekToEnd:最大オフセットからデータを消費することを示します。これにより、空のデータが読み取られます。

このパラメーターまたは beginDateTime のいずれかを指定する必要があります。

endOffset

データ消費の終了オフセットです。これは、データ消費タスクがいつ終了するかを制御するために使用されます。

このパラメーターまたは endDateTime のいずれかを指定する必要があります。

skipExceedRecord

Kafka は public ConsumerRecords<K, V> poll(final Duration timeout) を使用してデータを消費します。1 回の poll 呼び出しで、指定された endOffset または endDateTime を超えるデータがフェッチされる場合があります。このパラメーターは、この超過データを送信先に書き込むかどうかを決定します。タスクは自動オフセットコミットを使用するため、以下を推奨します:

  • 0.10.2 より前の Kafka バージョンの場合:skipExceedRecord を false に設定します。

  • Kafka 0.10.2 以降の場合:skipExceedRecord を true に設定します。

いいえ。デフォルト値は false です。

partition

Kafka トピックには複数のパーティション (partition) があります。デフォルトでは、データ同期タスクはトピック内のすべてのパーティションをカバーするオフセット範囲からデータを読み取ります。partition を指定して、単一のパーティションのオフセット範囲からのみデータを読み取ることもできます。

いいえ。デフォルト値はありません。

kafkaConfig

データ消費のために KafkaConsumer クライアントを作成する際に、bootstrap.serversauto.commit.interval.mssession.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 です。

  • このパラメーターが true に設定され、コンシューマーが poll リクエストからデータを受信しなかった場合、トピックパーティションの最大オフセットに到達したかどうかを確認します。すべてのパーティションで最大オフセットに到達した場合、タスクはすぐに停止します。それ以外の場合は、トピックからデータをプルし続けます。

  • このパラメーターが false に設定され、コンシューマーが poll リクエストからデータを受信しなかった場合、チェックを実行せずにタスクをすぐに停止します。

いいえ。デフォルト値は 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 です。

  • text に設定すると、すべての列が fieldDelimiter で指定された区切り文字を使用して連結されます。

  • JSON に設定すると、すべての列が column パラメーターで指定されたフィールド名に基づいて JSON 文字列に連結されます。

たとえば、ソースレコードに 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

データを書き込む送信先テーブルのフィールドをカンマで区切って指定します。例:"column": ["id", "name", "age"]

valueIndex が設定されておらず、writeMode が JSON に設定されている場合、このパラメーターはソースレコードの列値の JSON 構造内のフィールド名を定義します。例:"column": [{"name":id","type":"JSON_NUMBER"}, {"name":"name","type":"JSON_STRING"}, {"name":"age","type":"JSON_NUMBER"}]

  • ソースレコードの列数が column で設定されたフィールド名の数より多い場合、書き込み中にデータは切り捨てられます。例:

    ソースレコードに a、b、c の値を持つ 3 つの列があり、column が [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"}] として設定されている場合、書き込まれる Kafka レコードの値は文字列 {"col1":"a","col2":"b"} です。

  • ソースレコードの列数が column で指定されたフィールド数より少ない場合、余分なフィールドは null または nullValueFormat で指定された文字列で埋められます。例:

    ソースレコードに a と b の値を持つ 2 つの列があり、column が [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"},{"name":"col3","type":"JSON_STRING"}] として設定されている場合、書き込まれる Kafka レコードの値は文字列 {"col1":"a","col2":"b","col3":null} です。valueIndex が設定されているか、writeMode が text に設定されている場合、このパラメーターは無効です。

  • JSON フィールドタイプが設定されていない場合、デフォルトのフィールドタイプは JSON_STRING です。

  • 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 の有効な値は次のとおりです:

  • 0:書き込み成功の確認なし。

  • 1:プライマリレプリカへの書き込み成功の確認。

  • all:すべてのレプリカへの書き込み成功の確認。

いいえ

付録: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 フィールドに書き込みます。たとえば、ソースレコードの列値が整数 123 で、column[{"name":"col1","type":"JSON_STRING"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":"123"} です。

JSON_NUMBER

ソースレコードの列値を数値に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 1.23 で、column[{"name":"col1","type":"JSON_NUMBER"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":1.23} です。

JSON_BOOL

ソースレコードの列値をブール値に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 true で、column が [{"name":"col1","type":"JSON_BOOL"} として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":true} です。

JSON_ARRAY

ソースレコードの列値を JSON 配列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 [1,2,3] で、column[{"name":"col1","type":"JSON_ARRAY"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":[1,2,3]} です。

JSON_MAP

ソースレコードの列値を JSON オブジェクトに変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が文字列 {"k1":"v1"} で、column[{"name":"col1","type":"JSON_MAP"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":{"k1":"v1"}} です。

JSON_BASE64

ソース列のバイト配列を BASE64 エンコードされた文字列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が 16 進数で 0x01 0x02 と表される 2 バイト配列で、column[{"name":"col1","type":"JSON_BASE64"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":"AQI="} です。

JSON_HEX

ソース列のバイト配列を 16 進文字列に変換し、JSON フィールドに書き込みます。たとえば、ソースレコードの列値が 16 進数で 0x01 0x02 と表される 2 バイト配列で、column[{"name":"col1","type":"JSON_HEX"}] として設定されている場合、Kafka レコードに書き込まれる値は文字列 {"col1":"0102"} です。