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

Realtime Compute for Apache Flink:Kafka YAML コネクタ

最終更新日:Sep 16, 2026

Kafka コネクタは、Flink CDC データインジェストジョブでソースまたはシンクとして使用できます。このトピックでは、Kafka YAML コネクタの構文、パラメーター、および使用例について説明します。

前提条件

次のいずれかの方法でご利用のクラスターに接続します。

  • ApsaraMQ for Kafka クラスターへの接続

    • Kafka クラスターのバージョンが 0.11 以降であること。

    • ApsaraMQ for Kafka クラスターが作成されていること。詳細については、「リソースの作成」をご参照ください。

    • Flink ワークスペースと Kafka クラスターが同じ VPC 内にあり、ApsaraMQ for Kafka が Flink をホワイトリストに追加していること。詳細については、「ホワイトリストの設定」をご参照ください。

    重要

    ApsaraMQ for Kafka への書き込みに関する制限事項:

    • ApsaraMQ for Kafka は、zstd 圧縮フォーマットでのデータ書き込みをサポートしていません。

    • ApsaraMQ for Kafka は、べき等書き込みやトランザクション書き込みをサポートしていないため、Kafka シンクテーブルが提供する 1 回限りのセマンティクスは利用できません。VVR 8.0.0 以降、Kafka コネクタが使用するオープンソースの Kafka クライアントはバージョン 3.x にアップグレードされ、properties.enable.idempotence プロパティのデフォルト値が true になり、べき等書き込みが明示的に有効になります。したがって、VVR 8.0.0 以降で ApsaraMQ for Kafka に書き込む場合は、シンクテーブルに properties.enable.idempotence=false を明示的に追加して、べき等書き込みを無効にし、書き込みの失敗を回避する必要があります。ApsaraMQ for Kafka のストレージエンジンの比較と機能制限については、「ストレージエンジンの比較」をご参照ください。

  • セルフマネージド Apache Kafka クラスターへの接続

    • セルフマネージドの Apache Kafka クラスターでは、バージョン 0.11 以降が実行されています。

    • Flink とセルフマネージド Apache Kafka クラスター間のネットワークが接続されていること。インターネット経由でセルフマネージドクラスターに接続するには、「ネットワーク接続」をご参照ください。

    • Apache Kafka 2.8 のクライアント設定項目のみがサポートされています。詳細については、Apache Kafka のコンシューマーおよびプロデューサーの設定に関するドキュメントをご参照ください。

制限事項

  • VVR 11.1 以降では、Flink CDC データインジェストのデータソースとして Kafka を使用することを推奨します。

  • JSON、Debezium JSON、Canal JSON フォーマットのみがサポートされています。その他のデータ形式はサポートされていません。

  • ソースの場合、複数のパーティションに分散された同一テーブルのデータは、VVR 8.0.11 以降でのみサポートされます。

注意事項

現在、Flink および Kafka コミュニティの設計上の制限により、トランザクション書き込みは推奨されません。sink.delivery-guarantee = exactly-once を設定すると、Kafka コネクタはトランザクション書き込みを有効にします。以下の 3 つの問題が知られています。

  • 各チェックポイントでトランザクション ID が生成されます。チェックポイントの間隔が短すぎると、生成されるトランザクション ID が過多になります。Kafka クラスターのコーディネーターでメモリが不足し、クラスターの安定性が損なわれる可能性があります。

  • 各トランザクションでプロデューサーインスタンスが作成されます。同時にコミットされるトランザクションが多すぎると、TaskManager でメモリが不足し、Flink ジョブの安定性が損なわれる可能性があります。

  • 複数の Flink ジョブが同じ sink.transactional-id-prefix を使用すると、生成されるトランザクション ID が競合する可能性があります。1 つのジョブがデータ書き込みに失敗すると、Kafka パーティションの LSO (Log Start Offset) の進行が妨げられ、そのパーティションからデータを読み取るすべてのコンシューマーに影響します。

1 回限りのセマンティクスが必要な場合は、Upsert Kafka コネクタを使用してプライマリキーテーブルにデータを書き込み、プライマリキーに依存してべき等性を確保してください。トランザクション書き込みを使用する必要がある場合は、「1 回限りのセマンティクス」をご参照ください。

構文

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: ${kafka.topic}
sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: localhost:9092

設定項目

  • 全般

    パラメーター

    説明

    必須

    データの型

    デフォルト値

    注

    type

    ソースまたはシンクのタイプ。

    はい

    String

    なし

    値は kafka である必要があります。

    name

    ソースまたはシンクの名前。

    いいえ

    String

    なし

    なし

    properties.bootstrap.servers

    Kafka ブローカーのアドレス。

    はい

    String

    なし

    フォーマットは host:port,host:port,host:port です。複数のアドレスはカンマ (,) で区切ります。

    properties.*

    Kafka クライアントに直接渡される設定。

    いいえ

    String

    なし

    サフィックスは、Apache Kafka の公式ドキュメントで定義されているプロデューサーまたはコンシューマーの設定項目である必要があります。

    Flink は properties. プレフィックスを削除し、残りの設定を Kafka クライアントに渡します。たとえば、'properties.allow.auto.create.topics' = 'false' を使用して、トピックの自動作成を無効にできます。

    key.format

    Kafka メッセージのキー部分を読み書きするために使用されるフォーマット。

    いいえ

    String

    なし

    • ソースの場合、json のみがサポートされます。

    • シンクの場合、有効な値:

      • csv

      • json

    説明

    このパラメーターは VVR 11.0.0 以降でのみサポートされます。

    value.format

    Kafka メッセージの値部分を読み書きするために使用されるフォーマット。

    いいえ

    String

    debezium-json

    • ソースの場合、有効な値:

      • debezium-json 

      • canal-json

      • json

    • シンクの場合、有効な値:

      • debezium-json 

      • canal-json

      • canal-protobuf

    説明
    • debezium-json および canal-json フォーマットは、VVR 8.0.10 以降でのみサポートされます。

    • json フォーマットは、VVR 11.0.0 以降でのみサポートされます。

  • ソース

    パラメーター

    説明

    必須

    データの型

    デフォルト値

    注

    topic

    読み取るトピックの名前。

    いいえ

    String

    なし

    複数のトピック名はセミコロン (;) で区切ります (例:topic-1;topic-2)。

    説明

    topic と topic-pattern のいずれか一方のみを指定できます。

    topic-pattern

    読み取るトピックの名前に一致する正規表現。ジョブの実行時に、正規表現に一致するすべてのトピックが読み取られます。

    いいえ

    String

    なし

    例:

    • user_event_.*:名前が user_event_ で始まるすべてのトピックに一致します。

    • prod\.logs\..*:名前が prod.logs. プレフィックスを持つトピックに一致します (. はエスケープする必要があります)。

    説明

    topic と topic-pattern のいずれか一方のみを指定できます。

    properties.group.id

    コンシューマーグループの ID。

    いいえ

    String

    なし

    指定されたグループ ID を初めて使用する場合は、properties.auto.offset.reset を earliest または latest に設定して、初期起動オフセットを指定する必要があります。

    scan.startup.mode

    Kafka がデータを読み取る開始オフセット。

    いいえ

    String

    group-offsets

    有効な値:

    • earliest-offset:Kafka パーティションの最小オフセットからデータを読み取ります。

    • latest-offset:Kafka パーティションの最大オフセットからデータを読み取ります。

    • group-offsets (デフォルト):properties.group.id で指定されたコンシューマーグループのコミット済みオフセットからデータを読み取ります。

    • timestamp:scan.startup.timestamp-millis で指定されたタイムスタンプからデータを読み取ります。

    • specific-offsets:scan.startup.specific-offsets で指定されたオフセットからデータを読み取ります。

    説明

    このパラメーターは、ジョブがステートなしで開始された場合にのみ有効です。ジョブがチェックポイントから再開されたり、ステートから回復されたりすると、読み取りはステートに保存された進行状況から再開されます。

    scan.startup.specific-offsets

    起動モードが specific-offsets の場合の各パーティションの起動オフセット。

    いいえ

    String

    なし

    例:partition:0,offset:42;partition:1,offset:300

    scan.startup.timestamp-millis

    起動モードが timestamp の場合の起動タイムスタンプ。

    いいえ

    Long

    なし

    単位:ミリ秒。

    scan.topic-partition-discovery.interval

    Kafka のトピックとパーティションを動的に検出する間隔。

    いいえ

    Duration

    5 分

    デフォルトのパーティション検出間隔は 5 分です。この機能を無効にするには、パーティション検出間隔を明示的に 0 以下の値に設定します。動的パーティション検出が有効な場合、Kafka ソースは新しいパーティションを自動的に検出し、そこからデータを読み取ります。topic-pattern モードでは、ソースは既存のトピックの新しいパーティションからデータを読み取り、正規表現に一致する新しいトピックのすべてのパーティションからデータを読み取ります。

    scan.check.duplicated.group.id

    properties.group.id で指定されたコンシューマーグループが重複しているかどうかを確認するかどうかを指定します。

    いいえ

    Boolean

    false

    有効な値:

    • true:ジョブが開始される前にコンシューマーグループの競合をチェックします。競合が存在する場合、ジョブはエラーを報告します。これにより、既存のコンシューマーグループとの競合を回避できます。

    • false:コンシューマーグループの競合をチェックせずにジョブを直接開始します。

    schema.inference.strategy

    スキーマ解析ポリシー。

    いいえ

    String

    continuous

    有効な値:

    • continuous:すべてのレコードのスキーマを解析します。2 つの連続するスキーマに互換性がない場合、より広範なスキーマが解析され、スキーマ変更イベントが生成されます。

    • static:ジョブの開始時にスキーマを 1 回だけ解析します。その後、データは初期スキーマに基づいて解析され、スキーマ変更イベントは生成されません。

    説明

    scan.max.pre.fetch.records

    初期スキーマ解析中に各パーティションから消費および解析するメッセージの最大数。

    いいえ

    Int

    50

    ジョブがデータを読み取って処理する前に、指定された数の最新メッセージが各パーティションから事前に消費され、スキーマ情報が初期化されます。

    key.fields-prefix

    メッセージキーから解析されたフィールド名に追加されるカスタムプレフィックス。これにより、Kafka メッセージキーが解析された後の命名競合を回避できます。

    いいえ

    String

    なし

    たとえば、このパラメーターを key_ に設定し、キーに a という名前のフィールドが含まれている場合、キーが解析された後、フィールド名は key_a になります。

    説明

    key.fields-prefix の値は、value.fields-prefix のプレフィックスにすることはできません。

    value.fields-prefix

    メッセージ値から解析されたフィールド名に追加されるカスタムプレフィックス。これにより、Kafka メッセージ値が解析された後の命名競合を回避できます。

    いいえ

    String

    なし

    たとえば、このパラメーターを value_ に設定し、値に b という名前のフィールドが含まれている場合、値が解析された後、フィールド名は value_b になります。

    説明

    value.fields-prefix の値は、key.fields-prefix のプレフィックスにすることはできません。

    metadata.list

    ダウンストリームに渡すメタデータ列。

    いいえ

    String

    なし

    利用可能なメタデータ列は topic、partition、offset、timestamp、timestamp-type、headers、leader-epoch、__raw_key__、および __raw_value__ です。複数の列はカンマ (,) で区切ります。

    説明

    __raw_key__ および __raw_value__ メタデータ列は、VVR 11.6 以降でのみ利用可能です。

    scan.value.initial-schemas.ddls

    DDL 文を使用して指定された、特定のテーブルの初期スキーマ。

    いいえ

    String

    なし

    複数の DDL 文は ; で区切ります。たとえば、CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT); を使用して、db1.t1 および db1.t2 テーブルの初期スキーマを指定します。

    DDL 文のテーブルスキーマは、送信先テーブルと一致し、Flink SQL 構文規則に準拠している必要があります。

    説明

    このパラメーターは VVR 11.5 以降でのみサポートされます。

    ingestion.ignore-errors

    データ解析中のエラーを無視するかどうかを指定します。

    いいえ

    Boolean

    false

    説明

    このパラメーターは VVR 11.5 以降でのみサポートされます。

    ingestion.error-tolerance.max-count

    データ解析エラーが無視される場合に、ジョブが失敗するまでの累積解析エラー数。

    いいえ

    Integer

    -1

    このパラメーターは、ingestion.ignore-errors が有効な場合にのみ有効です。デフォルト値 -1 は、解析例外がジョブの失敗を引き起こさないことを示します。

    説明

    このパラメーターは VVR 11.5 以降でのみサポートされます。

    scan.duplicate-field.strategy

    キー部分と値部分から解析された重複フィールド名の処理方法を指定します。

    いいえ

    String

    EXCEPTION

    有効な値:

    • EXCEPTION:キーと値に重複フィールドが存在する場合に例外をスローします。これは VVR 11.6 以前のデフォルトの動作です。

    • PREFER_KEY:フィールドが重複している場合、キーフィールドの値を優先します。

    • PREFER_VALUE:フィールドが重複している場合、値フィールドの値を優先します。

    説明

    このパラメーターは VVR 11.7 以降でのみサポートされます。

    • Debezium JSON フォーマットのソーステーブル

      パラメーター

      必須

      データの型

      デフォルト値

      説明

      debezium-json.distributed-tables

      いいえ

      Boolean

      false

      Debezium JSON の単一テーブルのデータが複数のパーティションに存在する場合は、このオプションを有効にします。

      説明

      このパラメーターは VVR 8.0.11 以降でのみサポートされます。

      重要

      このパラメーターを変更した後は、ステートなしでジョブを開始する必要があります。

      debezium-json.schema-include

      いいえ

      Boolean

      false

      Debezium Kafka Connect を設定する際、Kafka 設定値 value.converter.schemas.enable を有効にして、メッセージにスキーマを含めることができます。このオプションは、Debezium JSON メッセージにスキーマが含まれるかどうかを指定します。

      有効な値:

      • true:Debezium JSON メッセージにスキーマが含まれます。

      • false:Debezium JSON メッセージにスキーマが含まれません。

      debezium-json.ignore-parse-errors

      いいえ

      Boolean

      false

      有効な値:

      • true:解析例外が発生した場合、現在の行をスキップします。

      • false (デフォルト):エラーを報告し、ジョブの開始に失敗します。

      debezium-json.infer-schema.primitive-as-string

      いいえ

      Boolean

      false

      テーブルスキーマの解析時に、すべての型を String として解析するかどうかを指定します。

      有効な値:

      • true:すべてのプリミティブ型を String として解析します。

      • false (デフォルト):基本ルールに基づいてデータを解析します。

      debezium-json.infer-schema.string-type-inference

      いいえ

      Boolean

      true

      文字列フィールドを TIME、DATE、または TIMESTAMP 型として推論しようとするかどうかを指定します。このパラメーターを false に設定すると、推論はスキップされ、フィールドは STRING のままになります。

      説明

      このパラメーターは VVR 11.8 以降でのみサポートされます。

    • Canal JSON フォーマットのソーステーブル

      パラメーター

      必須

      データの型

      デフォルト値

      説明

      canal-json.distributed-tables

      いいえ

      Boolean

      false

      Canal JSON の単一テーブルのデータが複数のパーティションに存在する場合は、このオプションを有効にします。

      説明

      このパラメーターは VVR 8.0.11 以降でのみサポートされます。

      重要

      このパラメーターを変更した後は、ステートなしでジョブを開始する必要があります。

      canal-json.database.include

      いいえ

      String

      なし

      Canal レコードのデータベースメタデータフィールドに一致するオプションの正規表現。これにより、指定されたデータベースの変更ログレコードのみが読み取られます。正規表現文字列は Java Pattern と互換性があります。

      canal-json.table.include

      いいえ

      String

      なし

      Canal レコードのテーブルメタデータフィールドに一致するオプションの正規表現。これにより、指定されたテーブルの変更ログレコードのみが読み取られます。正規表現文字列は Java Pattern と互換性があります。

      canal-json.ignore-parse-errors

      いいえ

      Boolean

      false

      有効な値:

      • true:解析例外が発生した場合、現在の行をスキップします。

      • false (デフォルト):エラーを報告し、ジョブの開始に失敗します。

      canal-json.infer-schema.primitive-as-string

      いいえ

      Boolean

      false

      テーブルスキーマの解析時に、すべての型を String として解析するかどうかを指定します。

      有効な値:

      • true:すべてのプリミティブ型を String として解析します。

      • false (デフォルト):基本ルールに基づいてデータを解析します。

      canal-json.infer-schema.strategy

      いいえ

      String

      AUTO

      テーブルスキーマの解析時に使用される解析ポリシー。

      有効な値:

      • AUTO (デフォルト):JSON データを解析して型を自動的に解析します。データに sqlType フィールドが含まれていない場合は、解析の失敗を避けるために AUTO を使用することを推奨します。

      • SQL_TYPE:Canal JSON データの sqlType 配列に基づいて型を解析します。データに sqlType フィールドが含まれている場合は、より正確な型を得るために canal-json.infer-schema.strategy を SQL_TYPE に設定することを推奨します。

      • MYSQL_TYPE:Canal JSON データの mysqlType 配列に基づいて型を解析します。

      Kafka の Canal JSON データに sqlType フィールドが含まれており、より正確な型マッピングが必要な場合は、canal-json.infer-schema.strategy を SQL_TYPE に設定することを推奨します。

      sqlType の型マッピングルールについては、「Canal JSON スキーマ解析」をご参照ください。

      説明
      • このパラメーターは VVR 11.1 以降でのみサポートされます。

      • MYSQL_TYPE は VVR 11.3 以降でのみサポートされます。

      canal-json.mysql.treat-mysql-timestamp-as-datetime-enabled

      いいえ

      Boolean

      true

      MySQL の timestamp 型を CDC の timestamp 型にマッピングするかどうかを指定します。

      • true (デフォルト):MySQL の timestamp 型を CDC の timestamp 型にマッピングします。

      • false:MySQL の timestamp 型を CDC の timestamp_ltz 型にマッピングします。

      canal-json.mysql.treat-tinyint1-as-boolean.enabled

      いいえ

      Boolean

      true

      MYSQL_TYPE 解析を使用する場合に、MySQL の tinyint(1) 型を CDC の boolean 型にマッピングするかどうかを指定します。

      • true (デフォルト):MySQL の tinyint(1) 型を CDC の boolean 型にマッピングします。

      • false:MySQL の tinyint(1) 型を CDC の tinyint(1) 型にマッピングします。

      このパラメーターは、canal-json.infer-schema.strategy が MYSQL_TYPE に設定されている場合にのみ有効です。

      canal-json.infer-schema.string-type-inference

      いいえ

      Boolean

      true

      文字列フィールドを TIME、DATE、または TIMESTAMP 型として推論しようとするかどうかを指定します。このパラメーターを false に設定すると、推論はスキップされ、フィールドは STRING のままになります。

      説明

      このパラメーターは VVR 11.8 以降でのみサポートされます。

    • JSON フォーマットのソーステーブル

      パラメーター

      必須

      データの型

      デフォルト値

      説明

      json.timestamp-format.standard

      いいえ

      String

      SQL

      入出力のタイムスタンプフォーマット。有効な値:

      • SQL:入力タイムスタンプを yyyy-MM-dd HH:mm:ss.s{precision} フォーマットで解析します (例:2020-12-30 12:13:14.123)。

      • ISO-8601:入力タイムスタンプを yyyy-MM-ddTHH:mm:ss.s{precision} フォーマットで解析します (例:2020-12-30T12:13:14.123)。

      json.ignore-parse-errors

      いいえ

      Boolean

      false

      有効な値:

      • true:解析例外が発生した場合、現在の行をスキップします。

      • false (デフォルト):エラーを報告し、ジョブの開始に失敗します。

      json.infer-schema.primitive-as-string

      いいえ

      Boolean

      false

      テーブルスキーマの解析時に、すべての型を String として解析するかどうかを指定します。

      有効な値:

      • true:すべてのプリミティブ型を String として解析します。

      • false (デフォルト):基本ルールに基づいてデータを解析します。

      json.infer-schema.flatten-nested-columns.enable

      いいえ

      Boolean

      false

      解析時に JSON データのネストされた列を再帰的に展開するかどうかを指定します。有効な値:

      • true:ネストされた列を再帰的に展開します。

      • false (デフォルト):ネストされた列を String として扱います。

      json.decode.parser-table-id.fields

      いいえ

      String

      なし

      JSON データの解析時に、特定の JSON フィールドの値を使用して tableId を生成するかどうかを指定します。複数のフィールドは , で区切ります。たとえば、JSON データが {"col0":"a", "col1","b", "col2","c"} の場合、生成される結果は次のようになります。

      設定

      tableId

      col0

      a

      col0,col1

      a.b

      col0,col1,col2

      a.b.c

      json.infer-schema.fixed-types

      いいえ

      String

      なし

      JSON データの解析時に特定のフィールドの型を指定します。複数のフィールドは , で区切ります。たとえば、id BIGINT, name VARCHAR(10) は、JSON データの id フィールドの型を BIGINT に、name フィールドの型を VARCHAR(10) に指定します。

      説明
      • このパラメーターは VVR 11.5 以降でのみサポートされます。

      • VVR 11.5 でこのパラメーターを使用する場合は、scan.max.pre.fetch.records: 0 も追加する必要があります。

      json.decode.converter-class

      いいえ

      String

      なし

      実装クラスの完全修飾名。コンバーターは、JSON データが解析される前に JSON byte[] を変更できます。

      説明

      このパラメーターは VVR 11.6 以降でのみサポートされます。

      json.decode.empty-value-as-delete.enabled

      いいえ

      Boolean

      false

      Kafka の compacted トピック内のトゥームストーンメッセージ (空の値を持つ) を DELETE イベントとして解析するかどうかを指定します。これは、空の値が削除を示すシナリオ (compacted トピックのミラーリングや CDC の削除シグナルなど) に適用されます。

      説明

      このパラメーターは VVR 11.7 以降でのみサポートされます。

      json.infer-schema.string-type-inference

      いいえ

      Boolean

      true

      文字列フィールドを TIME、DATE、または TIMESTAMP 型として推論しようとするかどうかを指定します。このパラメーターを false に設定すると、推論はスキップされ、フィールドは STRING のままになります。

      説明

      このパラメーターは VVR 11.8 以降でのみサポートされます。

  • シンク

    パラメーター

    説明

    必須

    データの型

    デフォルト値

    注

    type

    シンクのタイプ。

    はい

    String

    なし

    値は kafka である必要があります。

    name

    シンクの名前。

    いいえ

    String

    なし

    なし

    topic

    Kafka トピックの名前。

    いいえ

    String

    なし

    このオプションを有効にすると、すべてのデータがこのトピックに書き込まれます。

    説明

    このオプションを無効にすると、各レコードは TableID 文字列 (. で結合して生成) の名前を持つトピックに書き込まれます (例:databaseName.tableName)。

    partition.strategy

    Kafka パーティションにデータを書き込むために使用されるポリシー。

    いいえ

    String

    all-to-zero

    有効な値:

    • all-to-zero (デフォルト):すべてのデータをパーティション 0 に書き込みます。

    • hash-by-key:プライマリキーのハッシュ値に基づいて複数のパーティションにデータを書き込みます。これにより、同じプライマリキーを持つデータが同じパーティションに順番に書き込まれることが保証されます。

    sink.tableId-to-topic.mapping

    ソーステーブル名と送信先 Kafka トピック名の間のマッピング。

    いいえ

    String

    なし

    各マッピングは ; で区切ります。ソーステーブル名と送信先 Kafka トピック名は : で区切ります。テーブル名は正規表現にすることができます。同じトピックにマッピングされる複数のテーブルは , で結合できます。例:mydb.mytable1:topic1;mydb.mytable2:topic2。

    説明

    このパラメーターを使用すると、元のテーブル名情報を保持したまま、マッピングされたトピックを変更できます。

    • Canal JSON フォーマットのシンクテーブル

      パラメーター

      必須

      データの型

      デフォルト値

      説明

      canal-json.serialize.update.keep-changed-fields-only

      いいえ

      Boolean

      false

      canal-json フォーマットの UPDATE メッセージの old 部分に、変更されたフィールドの変更前の値のみを含めるかどうかを指定します。

      説明

      このパラメーターは VVR 11.8 以降でのみサポートされます。

    • Debezium JSON フォーマットのシンクテーブル

      パラメーター

      必須

      データの型

      デフォルト値

      説明

      debezium-json.include-schema.enabled

      いいえ

      Boolean

      false

      Debezium JSON データにスキーマ情報を含めるかどうかを指定します。

      debezium-json.emit.full-table-id.enabled

      いいえ

      Boolean

      false

      完全な 3 部構成のテーブル ID を Debezium JSON メタデータフィールドに書き込むかどうかを指定します。

      このパラメーターが有効な場合のマッピング:

      CDC テーブル ID の部分

      Debezium JSON キー

      Namespace

      db

      Schema

      schema

      Table

      table

      このパラメーターが無効な場合のマッピング:

      CDC テーブル ID の部分

      Debezium JSON キー

      Namespace

      なし

      Schema

      db

      Table

      table

      説明

      このパラメーターは VVR 11.6 以降でのみサポートされます。

既存のカタログの再利用

VVR 11.5 以降では、Flink CDC データインジェストジョブで [カタログ] ページで作成された組み込みの Kafka カタログを直接参照できます。これにより、接続プロパティを手動で記述する必要がなくなります。

source:
  type: kafka
  using.built-in-catalog: kafka_catalog

現在、データインジェストジョブは、次の Kafka カタログパラメーターの自動再利用をサポートしています。

  • properties.bootstrap.servers

  • format

  • key.fields-prefix

  • value.fields-prefix

  • timestamp-format.standard

  • infer-schema.flatten-nested-columns.enable

  • infer-schema.primitive-as-string

  • max.fetch.records

自動的に再利用されるパラメーターをオーバーライドするには、対応する YAML パラメーターを明示的に指定します。YAML パラメーターが優先されます。

設定例

  • データインジェストジョブのソースとして Kafka を使用する:

    source:
      type: kafka
      name: Kafka source
      properties.bootstrap.servers: ${kafka.bootstraps.server}
      topic: ${kafka.topic}
      value.format: ${value.format}
      scan.startup.mode: ${scan.startup.mode}
     
    sink:
      type: hologres
      name: Hologres sink
      endpoint: <yourEndpoint>
      dbname: <yourDbname>
      username: ${secret_values.ak_id}
      password: ${secret_values.ak_secret}
      sink.type-normalize-strategy: BROADEN
  • データインジェストジョブのシンクとして Kafka を使用する:

    source:
      type: mysql
      name: MySQL Source
      hostname: ${secret_values.mysql.hostname}
      port: ${mysql.port}
      username: ${secret_values.mysql.username}
      password: ${secret_values.mysql.password}
      tables: ${mysql.source.table}
      server-id: 8601-8604
    
    sink:
      type: kafka
      name: Kafka Sink
      properties.bootstrap.servers: ${kafka.bootstraps.server}
    
    route:
      - source-table: ${mysql.source.table}
        sink-table: ${kafka.topic}

    この例では、route モジュールを使用して、ソーステーブルが書き込まれる Kafka トピックの名前を設定します。

説明

ApsaraMQ for Kafka は、デフォルトではトピックの自動作成を有効にしていません。詳細については、「トピックの自動作成に関するよくある質問」をご参照ください。ApsaraMQ for Kafka にデータを書き込む前に、対応するトピックを作成する必要があります。詳細については、「ステップ 3:リソースの作成」をご参照ください。

使用例

以下に、典型的なシナリオの設定例を示します。

単一トピックの読み取り

次の例では、トピック customers を読み取り、データを Data Lake Formation に書き込みます。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true

この例では、JSON フォーマットで生成されるテーブル名は、デフォルトでトピック名と同じになります。

複数トピックの読み取り

次の例では、正規表現に一致する複数のトピックを読み取り、データを StarRocks コネクタに書き込みます。

source:
  type: kafka
  topic-pattern: user_event_.*
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: starrocks
  jdbc-url: jdbc:mysql://<yourFeHostname>:9030
  load-url: <yourFeHostname>:8030
  username: <yourUsername>
  password: ${secret_values.starrocks_password}
 
  # (任意) データ量が少ないジョブでは、データが長時間ディスクに書き込まれないのを避けるため、フラッシュ間隔を短くすることを推奨します (デフォルトは 300000、つまり 5 分)
  sink.buffer-flush.interval-ms: 5000
  # (任意) 上流が utf8mb4 文字セットの場合、テキストの切り捨てを避けるために 4 に設定することを推奨します (デフォルトは 3)
  unicode-char.max-bytes: 4
  # (任意) 自動テーブル作成のバケット数。StarRocks 2.5.7 以下では明示的に設定する必要がありますが、それ以上のバージョンでは自動的に推論できます
  table.create.num-buckets: 8
  # (任意) 自動テーブル作成のレプリケーション数。クラスターの状況に応じて設定します
  table.create.properties.replication_num: 3
  # (任意) StarRocks 3.2 以降では、テーブル構造の変更を高速化するために有効にすることを推奨します
  table.create.properties.fast_schema_evolution: true
  # 注意:transform を使用してプライマリキーを変更する場合、同時に sink.ignore.update-before: false を設定する必要があります。
  # そうしないと、古いプライマリキーに対応する行が下流に残ってしまいます

この例では、JSON フォーマットで生成されるテーブル名は、デフォルトでトピック名と同じになります。

キーの読み取りとフィールド競合の防止

キーと値の間のフィールド競合によるエラーを防ぐには、次のいずれかのソリューションを使用します。

  1. フィールドにプレフィックスを追加して、フィールド競合を回避します。

source:
  type: kafka
  topic: ${kafka.topic}
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  key.format: json
  value.format: json
  # key 部分のフィールド名に key_ プレフィックスを追加します
  key.fields-prefix: key_
  # value 部分のフィールド名に value_ プレフィックスを追加します
  value.fields-prefix: value_
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  1. 競合解決ポリシーを設定します。詳細については、「scan.duplicate-field.strategy」をご参照ください。次の設定では、キーのフィールドが優先され、重複が発生した場合に値の同じ名前のフィールドは無視されます。

source:
  type: kafka
  topic: ${kafka.topic}
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  key.format: json
  value.format: json
  # key のフィールドを優先し、同名の value のフィールドを無視します
  scan.duplicate-field.strategy: PREFER_KEY
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true

メタデータ列の追加

次の例では、トピック customers を読み取り、メタデータ列 topic と partition をフィールドに追加します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous
  metadata.list: topic,partition

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true

解析エラーの処理

データ解析エラーはジョブの失敗を引き起こします。解析エラーを許容するようにジョブを設定できます。これは一般的にダーティデータ収集と組み合わせて使用されます。

次の例では、解析エラーを完全に無視します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous
  # 解析エラーの無視を有効にします。デフォルトではすべての解析エラーを無視します
  ingestion.ignore-errors: true

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

次の例では、累積解析エラーが 30 回に達するとジョブが失敗します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous
  # 解析エラーの無視を有効にします
  ingestion.ignore-errors: true
  # 解析エラーが 30 回発生するとジョブが失敗します
  ingestion.error-tolerance.max-count: 30
  
sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

JSON データの読み取り

以下のセクションでは、JSON フォーマットのデータを読み取る一般的な方法について説明します。

解析用のテーブル ID の指定

デフォルトでは、JSON データのテーブル ID はトピック名です。データのフィールドの値をテーブル ID として指定できます。次の例では、db および tbl フィールドがテーブル ID として指定されています。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous
  # フィールド db と table を table id として指定します
  value.json.decode.parser-table-id.fields: db,tbl

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

フィールドタイプの指定

フィールドタイプはフィールド値から推論されます。一部の推論されたタイプは、期待するタイプではない場合があります。特定のフィールドに固定タイプを指定し、これらのフィールドタイプの後続の推論と進化をスキップできます。

次の設定では、4 つのフィールドのタイプを固定します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: json
  # フィールド db と tbl を table id として指定します
  value.json.decode.parser-table-id.fields: db,tbl
  # 特定のフィールドタイプを固定します
  value.json.infer-schema.fixed-types: 'db STRING, tbl STRING, id BIGINT, amount DECIMAL(18, 2)'
  # 宣言されていないフィールドの動的推論を許可します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Canal JSON データの読み取り

以下のセクションでは、Canal JSON フォーマットのデータを読み取る一般的な方法について説明します。

型推論ポリシー

デフォルトでは、Canal JSON データを読み取る際、データ値を使用してスキーマタイプを推論します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: canal-json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Canal JSON データに記録されているスキーマ情報 (sql type または mysql type) に基づいて型を推論することもできます。次の例では、mysql type を使用してスキーマを推論します。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: canal-json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous
  # mysql type 情報を使用してスキーマを推論します。SQL_TYPE に設定して sql type で推論することもできます
  value.canal-json.infer-schema.strategy: MYSQL_TYPE

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true
  
# ダーティデータ収集を有効にして、解析に失敗したデータを出力します
pipeline:
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Debezium JSON データの読み取り

次の例では、トピック customers を読み取り、データを Data Lake Formation に書き込みます。

source:
  type: kafka
  topic: customers
  properties.bootstrap.servers: localhost:9092
  properties.group.id: ${kafka.group.id}
  value.format: debezium-json
  # (任意) 各データのスキーマを動的に識別し、比較してスキーマ変更を生成します
  schema.inference.strategy: continuous

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (任意) コミットユーザー名。競合を避けるため、ジョブごとに異なるコミットユーザーを設定することを推奨します
  commit.user: your_job_name
  # (任意) 削除ベクトルを有効にして、読み取り性能を向上させます
  table.properties.deletion-vectors.enabled: true

生の MySQL バイナリログを Kafka に同期

Flink CDC データインジェストは、生の MySQL バイナリログコンテンツを Canal JSON に同期することをサポートしています。次のジョブは、複数のテーブルのバイナリログをトピック order_dw_tables に同期します。

source:
  type: mysql
  hostname: #{hostname}
  port: 3306
  username: #{username}
  password: #{password}
  tables: order_dw.\.*
  server-id: 28601-28604
  # (任意) 増分フェーズで新しく作成されたテーブルのデータを同期します
  scan.binlog.newly-added-table.enabled: true
  # (任意) テーブルコメントとフィールドコメントを同期します
  include-comments.enabled: true
  # (任意) TaskManager の OutOfMemory 問題を回避するために、無制限のチャンクを優先的に配信します
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (任意) 解析フィルタリングを有効にして、読み取りを高速化します
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # Canal JSON に mysqlType、sqlType、sql、isDdl などの情報を補足します
  include-binlog-meta.enable: true
  
sink:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: order_dw_tables
  # Kafka の value は Canal JSON changelog フォーマットを使用します
  value.format: canal-json
  # 日付型データをシリアル化する際に使用するフォーマットを指定します
  value.canal-json.timestamp-format.standard: SQL
  # バイナリログの順序を保証するために、データをパーティション 0 に一括で書き込みます
  partition.strategy: all-to-zero

高度な機能

スキーマ解析と変更同期ポリシー

Kafka コネクタは、現在既知のすべてのテーブルのスキーマを維持します。

テーブルスキーマ情報の初期化

テーブルスキーマ情報には、フィールドとデータ型情報、データベースとテーブル情報、およびプライマリキー情報が含まれます。これら 3 種類の情報は、次のように初期化されます。

  • フィールドとデータ型情報

データインジェストジョブは、データからテーブルのフィールドとデータ型を自動的に推論できます。ただし、一部のシナリオでは、特定のテーブルのフィールドと型を指定したい場合があります。フィールドタイプを指定する粒度に基づいて、テーブルスキーマ情報は次の 3 つのポリシーのいずれかを使用して初期化できます。

  1. プログラムによる完全な推論

Kafka データが読み取られる前に、Kafka コネクタは、各パーティションから scan.max.pre.fetch.records で指定されたメッセージ数までを事前に消費しようとし、各レコードのスキーマを解析し、これらのスキーマをマージしてテーブルスキーマ情報を初期化します。データが実際に消費される前に、初期化されたスキーマに基づいて対応するテーブル作成イベントが生成されます。

説明

Debezium JSON および Canal JSON フォーマットの場合、テーブル情報は個々のメッセージに含まれています。scan.max.pre.fetch.records で事前に消費されるメッセージには複数のテーブルのデータが含まれる可能性があるため、各テーブルで事前に消費されるレコード数を決定することはできません。事前消費とテーブルスキーマの初期化は、各パーティションのメッセージが実際に消費されて処理される前に 1 回だけ実行されます。後で新しいテーブルのデータが到着した場合、そのテーブルの最初のレコードから解析されたテーブルスキーマが初期スキーマとして使用され、そのテーブルに対して事前消費と初期化は再度実行されません。

重要

複数のパーティションに分散された単一テーブルのデータは、VVR 8.0.11 以降でのみサポートされます。このシナリオでは、debezium-json.distributed-tables または canal-json.distributed-tables を true に設定する必要があります。

  1. 初期テーブルスキーマの指定

一部のシナリオでは、たとえば、事前に作成されたダウンストリームテーブルに Kafka データを書き込む場合など、初期テーブルスキーマを自分で指定したい場合があります。この場合、scan.value.initial-schemas.ddls パラメーターを追加して初期テーブルスキーマを指定できます。次の例に設定を示します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # データ内の db、table フィールドを Table ID として使用します
  json.decode.parser-table-id.fields: db,table
  # 初期テーブル構造を設定します
  scan.value.initial-schemas.ddls: CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);

CREATE TABLE 文は、送信先テーブルのスキーマと一致している必要があります。この例では、db1.t1 テーブルの id フィールドの初期型は BIGINT、name フィールドの初期型は VARCHAR(10) に指定され、db1.t2 テーブルの id フィールドの初期型は BIGINT に指定されています。

CREATE TABLE 文は Flink SQL 構文を使用します。

  1. フィールドの固定タイプの指定

一部のシナリオでは、特定のフィールドのデータ型を固定したい場合があります。たとえば、TIMESTAMP 型として推論される可能性のある特定のフィールドを文字列として配信したい場合があります。この場合、json.infer-schema.fixed-types パラメーターを追加して初期テーブルスキーマを指定できます。これは、メッセージフォーマットが json の場合にのみ有効です。次の例に設定を示します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # 特定のフィールドを常に固定タイプに設定します
  json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
  scan.max.pre.fetch.records: 0

この例では、すべての id フィールドの型は BIGINT に、すべての name フィールドの型は VARCHAR(10) に固定されます。

ここで使用される型は、Flink SQL データ型と同じです。

  • データベースとテーブル情報

    • Canal JSON および Debezium JSON フォーマットの場合、テーブル情報は個々のメッセージから解析され、データベース名とテーブル名が含まれます。

    • JSON フォーマットの場合、テーブル情報にはデフォルトでテーブル名のみが含まれ、これはデータが存在するトピックの名前です。データにデータベースとテーブル情報が含まれている場合は、json.infer-schema.fixed-types を使用して、データベースとテーブル情報を含むフィールドを指定できます。これらのフィールドはデータベース名とテーブル名にマッピングされます。次の例に設定を示します。

      source:
        type: kafka
        name: Kafka Source
        properties.bootstrap.servers: host:9092
        topic: test-topic
        value.format: json
        scan.startup.mode: earliest-offset
        # col1 フィールドの値をデータベース名として、col2 フィールドの値をテーブル名として使用します
        json.decode.parser-table-id.fields: col1,col2

      この例では、各レコードは、データベース名が col1 フィールドの値で、テーブル名が col2 フィールドの値であるテーブルに配信されます。

  • プライマリキー情報

    • Canal JSON フォーマットの場合、テーブルのプライマリキーは JSON の pkNames フィールドに基づいて定義されます。

    • Debezium JSON および JSON フォーマットの場合、JSON にはプライマリキー情報が含まれていません。変換ルールを使用して、テーブルに手動でプライマリキーを追加できます。

      transform:
        - source-table: \.*.\.*
          projection: \*
          primary-keys: key1, key2

スキーマ解析とスキーマ変更

テーブルスキーマが初期化された後、schema.inference.strategy が static に設定されている場合、Kafka コネクタは初期テーブルスキーマに基づいて各メッセージの値を解析し、スキーマ変更イベントを生成しません。schema.inference.strategy が continuous に設定されている場合、Kafka コネクタは各 Kafka メッセージの値を解析し、メッセージの物理列を取得し、それらを現在維持されているスキーマと比較します。解析されたスキーマが現在のスキーマと一致しない場合、コネクタはスキーマをマージしようとし、対応するテーブルスキーマ変更イベントを生成します。マージルールは次のとおりです。

  • 解析された物理列に現在のスキーマに存在しないフィールドが含まれている場合、これらのフィールドはスキーマに追加され、add-nullable-column イベントが生成されます。

  • 解析された物理列に現在のスキーマに既に存在するフィールドが含まれていない場合、そのフィールドは保持され、その列のデータは NULL で埋められます。drop-column イベントは生成されません。

  • 両方に同じ名前の列が含まれている場合、列は次のように処理されます。

    • 型は同じだが精度が異なる場合、より高い精度の型が使用され、column-type-change イベントが生成されます。

    • 型が異なる場合、次のツリー構造の最も低い共通の親ノードが列の型として使用され、column-type-change イベントが生成されます。

      image

  • 現在サポートされているスキーマ変更ポリシーは次のとおりです。

    • 列の追加:対応する列が現在のスキーマの末尾に追加され、新しい列のデータが同期され、新しい列は null 許容列として設定されます。

    • 列の削除:drop-column イベントは生成されません。代わりに、列のデータは自動的に NULL 値で埋められます。

    • 列名の変更:これは列の追加と列の削除として扱われます。名前が変更された列は現在のスキーマの末尾に追加され、名前変更前の列のデータは NULL 値で埋められます。

    • 列の型の変更:

      • 列の型の変更をサポートするダウンストリームシステムの場合、データインジェストジョブは、ダウンストリームシンクが列の型の変更の処理をサポートした後、通常の列の型の変更をサポートします。たとえば、列を INT 型から BIGINT 型に変更できます。このような変更は、ダウンストリームシンクがサポートする列の型の変更ルールに依存します。異なるシンクテーブルは異なるルールをサポートします。サポートする列の型の変更ルールについては、シンクテーブルのドキュメントをご参照ください。

      • Hologres など、列の型の変更をサポートしないダウンストリームシステムの場合、型の拡張を使用できます。この場合、ジョブの開始時に、より広範な型を持つテーブルがダウンストリームシステムに作成されます。列の型の変更が発生した場合、ジョブはダウンストリームシンクが変更を受け入れられるかどうかをチェックし、列の型の変更に対する寛容なサポートを実装します。

  • サポートされていないスキーマ変更は次のとおりです。

    • プライマリキーやインデックスなどの制約の変更。

    • NOT NULL から NULLABLE への変更。

  • Canal JSON スキーマ解析

    Canal JSON データには、データ列の正確な型情報を記録するオプションの sqlType フィールドが含まれている場合があります。より正確なスキーマを取得するには、canal-json.infer-schema.strategy を SQL_TYPE に設定して、sqlType の型を使用できます。型マッピングは次のとおりです。

    JDBC 型

    型コード

    CDC 型

    BIT

    -7

    BOOLEAN

    BOOLEAN

    16

    TINYINT

    -6

    TINYINT

    SMALLINT

    -5

    SMALLINT

    INTEGER

    4

    INT

    BIGINT

    -5

    BIGINT

    DECIMAL

    3

    DECIMAL(38,18)

    NUMERIC

    2

    REAL

    7

    FLOAT

    FLOAT

    6

    DOUBLE

    8

    DOUBLE

    BINARY

    -2

    BYTES

    VARBINARY

    -3

    LONGVARBINARY

    -4

    BLOB

    2004

    DATE

    91

    DATE

    TIME

    92

    TIME

    TIMESTAMP

    93

    TIMESTAMP

    CHAR

    1

    STRING

    VARCHAR

    12

    LONGVARCHAR

    -1

    その他の型

ダーティデータの許容と収集

場合によっては、Kafka データソースに不正な形式のデータ (ダーティデータ) が含まれていることがあります。このダーティデータによって同期ジョブが頻繁に失敗して再起動するのを防ぐために、そのような異常なデータを無視するようにジョブを設定できます。次の例に設定を示します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # ダーティデータ許容機能を有効にします
  ingestion.ignore-errors: true
  # 1000 件のダーティデータを許容します
  ingestion.error-tolerance.max-count: 1000

この設定では、最大 1,000 件のダーティレコードが無視されるため、少量のダーティデータが存在する場合でもジョブは正常に実行されます。ダーティレコードの数がこのしきい値を超えると、ジョブは失敗し、データの検証を促します。

ダーティデータによってジョブが失敗しないようにするには、次の設定を使用します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # ダーティデータ許容機能を有効にします
  ingestion.ignore-errors: true
  # すべてのダーティデータを許容します
  ingestion.error-tolerance.max-count: -1

ダーティデータ許容ポリシーは、異常なデータによるジョブの頻繁な失敗を防ぎます。また、Kafka データプロデューサーの動作を調整するために、ダーティデータについて詳しく知りたい場合もあります。「ダーティデータ収集」で説明されているプロセスに従うことで、TaskManager ログでジョブのダーティデータを表示できます。次の例に設定を示します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: host:9092
  topic: test-topic
  value.format: json
  scan.startup.mode: earliest-offset
  # ダーティデータ許容機能を有効にします
  ingestion.ignore-errors: true
  # すべてのダーティデータを許容します
  ingestion.error-tolerance.max-count: -1

pipeline:
  dirty-data.collector:
    # ダーティデータを TaskManager のログファイルに書き込みます
    type: logger

テーブル名とトピック間のマッピングポリシー

データインジェストジョブのシンクとして Kafka を使用する場合、Kafka メッセージフォーマット (debezium-json または canal-json) にもテーブル名情報が含まれます。Kafka メッセージが後で消費されるとき、データ内のテーブル名は通常、トピック名ではなく実際のテーブル名として使用されます。したがって、テーブル名とトピック間のマッピングポリシーを慎重に設定する必要があります。

MySQL の 2 つのテーブル mydb.mytable1 と mydb.mytable2 を同期する必要があるとします。次の設定ポリシーが利用可能です。

1. マッピングポリシーを設定しない

マッピングポリシーがない場合、各テーブルはデータベース名.テーブル名形式で名前が付けられた対応するトピックに書き込まれます。したがって、mydb.mytable1 のデータは mydb.mytable1 という名前のトピックに書き込まれ、mydb.mytable2 のデータは mydb.mytable2 という名前のトピックに書き込まれます。次の例に設定を示します。

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}

2. マッピング用のルート規則を設定する (非推奨)

多くのシナリオでは、送信先トピックがデータベース名.テーブル名形式であることを望まず、指定されたトピックにデータを書き込みたい場合があります。この場合、マッピング用のルート規則を設定できます。次の例に設定を示します。

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  
 route:
  - source-table: mydb.mytable1,mydb.mytable2
    sink-table: mytable1

この場合、mydb.mytable1 と mydb.mytable2 からのすべてのデータは、単一のトピック mytable1 に書き込まれます。

ただし、ルート規則を使用して送信先トピック名を変更すると、Kafka メッセージ (debezium-json または canal-json 形式) のテーブル名情報も変更されます。この場合、Kafka メッセージ内のすべてのテーブル名は mytable1 になり、他のシステムがこのトピックの Kafka メッセージを消費するときに予期しない結果を引き起こす可能性があります。

3. sink.tableId-to-topic.mapping パラメーターをマッピング用に設定する (推奨)

ソーステーブル名情報を保持しながらテーブル名とトピック間のマッピング規則を設定するには、sink.tableId-to-topic.mapping パラメーターを使用できます。次の例に設定を示します。

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  sink.tableId-to-topic.mapping: mydb.mytable1,mydb.mytable2:mytable

または

source:
  type: mysql
  name: MySQL Source
  hostname: ${secret_values.mysql.hostname}
  port: ${mysql.port}
  username: ${secret_values.mysql.username}
  password: ${secret_values.mysql.password}
  tables: mydb.mytable1,mydb.mytable2
  server-id: 8601-8604

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: ${kafka.bootstraps.server}
  sink.tableId-to-topic.mapping: mydb.mytable1:mytable;mydb.mytable2:mytable

この場合、mydb.mytable1 と mydb.mytable2 からのすべてのデータは単一のトピック mytable1 に書き込まれ、Kafka メッセージ (debezium-json または canal-json 形式) のテーブル名情報は mydb.mytable1 または mydb.mytable2 のままです。他のシステムがこのトピックの Kafka メッセージを消費するとき、それらはソーステーブル名情報を正しく取得できます。

JSON コンバーターの実装

JSON データは常に完全に一貫した形式であるとは限らず、一部の処理要件は Transform モジュールでは満たせない場合があります。たとえば、2 つの列を新しい列にマージし、元の 2 つの列を削除し、スキーマ変更の同期が引き続き機能することを確認する必要がある場合があります。JSON データをより柔軟に処理するには、KafkaPayloadConverter インターフェイスを実装して、フレームワークが処理する前に JSON データを変更できます。この機能を使用するには、json.decode.converter-class 設定を追加し、実装クラスの完全修飾名に設定します。

JSON コンバーターを実装して使用するプロセスは次のとおりです。

  1. KafkaPayloadConverter インターフェイスを実装し、パッケージ化します。オープンソースのデモリポジトリには、既にいくつかの実装例が提供されています。

  2. パッケージ化されたファイルを追加の依存関係としてデータインジェストジョブに追加します。

  3. Source モジュールで、json.decode.converter-class 設定を指定します。次の例では、オープンソースプロジェクトの ArrayElementExtractorConverter クラスを使用します。

    source:                                                                                                                                                                                                                                                                                                            
      type: kafka                                                                                                                                                                                                                                                                                          
      properties.bootstrap.servers: localhost:9092                                                                                                                                                                                                                                                                     
      topic: my_cdc_topic                                                                                                                                                                                                                                                                                              
      properties.group.id: flink-cdc-group                                                                                                                                                                                                                                                                             
      scan.startup.mode: earliest-offset  
      value.format: json                                                                                                                                                                                                                                                                              
      # Kafka メッセージ value の JSON データにカスタムコンバーターを適用します
      value.json.decode.converter-class: org.apache.flink.cdc.connectors.kafka.ArrayElementExtractorConverter
  4. ジョブをデプロイして実行します。

key format と value format

VVR 11.5 以前では、キーフォーマットと値フォーマットの設定を区別できません。同じフォーマットを使用する場合、フォーマット設定はキーフォーマットと値フォーマットの両方に適用されます。

VVR 11.6 以降では、この動作が最適化されています。json フォーマットを例にとると、フォーマット設定項目は次のように渡されます。

  • フォーマットプレフィックス (json.infer-schema.primitive-as-string など):バージョンの互換性を維持するため、フォーマットプレフィックスを持つ設定は、デフォルトでキーフォーマットと値フォーマットの両方に適用されます。

  • キープレフィックス + フォーマットプレフィックス (key.json.infer-schema.primitive-as-string など):キープレフィックス + フォーマットプレフィックスを持つ設定は、キーフォーマットにのみ適用され、フォーマットプレフィックスを持つ設定よりも優先されます。

  • 値プレフィックス + フォーマットプレフィックス (value.json.infer-schema.primitive-as-string など):値プレフィックス + フォーマットプレフィックスを持つ設定は、値フォーマットにのみ適用され、フォーマットプレフィックスを持つ設定よりも優先されます。

次の Kafka ソース設定では、キーフォーマットと値フォーマットに異なる JSON コンバーターが設定されています。

source:                                                                                                                                                                                                                                                                                                            
  type: kafka                                                                                                                                                                                                                                                                                          
  properties.bootstrap.servers: localhost:9092                                                                                                                                                                                                                                                                     
  topic: my_cdc_topic                                                                                                                                                                                                                                                                                              
  properties.group.id: flink-cdc-group                                                                                                                                                                                                                                                                             
  scan.startup.mode: earliest-offset 
  key.format: json  
  value.format: json                                                                                                                                                                                                                                                                            
  # Kafka メッセージ key の JSON データにカスタムコンバーターを適用します
  key.json.decode.converter-class: com.example.KeyExampleConverter                                                                                                                                                                                                                                                                            
  # Kafka メッセージ value の JSON データにカスタムコンバーターを適用します
  value.json.decode.converter-class: com.example.ValueExampleConverter