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 クラスターへの接続
制限事項
-
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:300scan.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 回だけ解析します。その後、データは初期スキーマに基づいて解析され、スキーマ変更イベントは生成されません。
説明-
スキーマ解析の詳細については、「スキーマ解析と変更同期ポリシー」をご参照ください。
-
このパラメーターは VVR 8.0.11 以降でのみサポートされます。
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
dbSchema
schemaTable
tableこのパラメーターが無効な場合のマッピング:
CDC テーブル ID の部分
Debezium JSON キー
Namespace
なし
Schema
dbTable
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 フォーマットで生成されるテーブル名は、デフォルトでトピック名と同じになります。
キーの読み取りとフィールド競合の防止
キーと値の間のフィールド競合によるエラーを防ぐには、次のいずれかのソリューションを使用します。
-
フィールドにプレフィックスを追加して、フィールド競合を回避します。
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
-
競合解決ポリシーを設定します。詳細については、「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 つのポリシーのいずれかを使用して初期化できます。
-
プログラムによる完全な推論
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 に設定する必要があります。
-
初期テーブルスキーマの指定
一部のシナリオでは、たとえば、事前に作成されたダウンストリームテーブルに 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 構文を使用します。
-
フィールドの固定タイプの指定
一部のシナリオでは、特定のフィールドのデータ型を固定したい場合があります。たとえば、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 イベントが生成されます。

-
-
現在サポートされているスキーマ変更ポリシーは次のとおりです。
-
列の追加:対応する列が現在のスキーマの末尾に追加され、新しい列のデータが同期され、新しい列は 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 コンバーターを実装して使用するプロセスは次のとおりです。
-
KafkaPayloadConverter インターフェイスを実装し、パッケージ化します。オープンソースのデモリポジトリには、既にいくつかの実装例が提供されています。
-
パッケージ化されたファイルを追加の依存関係としてデータインジェストジョブに追加します。
-
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 -
ジョブをデプロイして実行します。
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