このトピックでは、Flink CDC データインジェストの YAML ジョブを使用して、ログデータを Alibaba Cloud Data Lake Formation (DLF) に書き込むためのベストプラクティスについて説明します。
データレイクへの Kafka ログのインジェスト
Flink CDC データインジェストとシンプルな YAML ジョブを使用して、ログデータをリアルタイムで迅速にデータレイクにインジェストします。システムは自動的にスキーマ推論を実行し、スキーマ進化をサポートします。
Apache Kafka の inventory トピックが、JSON フォーマットのログテーブルのデータを格納していると仮定します。次のサンプルジョブは、このデータを DLF の対応する結果テーブルに同期します。
source:
type: kafka
name: Kafka Source
# Kafka ブローカーのアドレス。
properties.bootstrap.servers: ${kafka.bootstrap.servers}
# 消費するトピック。
topic: inventory
# 最も早いオフセットから消費を開始するように指定します。
scan.startup.mode: earliest-offset
# Kafka メッセージ値のフォーマット。
value.format: json
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
json.infer-schema.flatten-nested-columns.enable: true
# (オプション) 最初の 100 件の解析例外をスキップします。例外の数が 100 を超えると、ジョブは失敗します。
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
# メタストアのタイプ。値は rest に固定されます。
catalog.properties.metastore: rest
# トークンプロバイダー。値は dlf に固定されます。
catalog.properties.token.provider: dlf
# DLF Rest Catalog Server にアクセスするための URI。フォーマットは http://[region-id]-vpc.dlf.aliyuncs.com です。例: http://cn-hangzhou-vpc.dlf.aliyuncs.com。
catalog.properties.uri: dlf_uri
# DLF カタログの名前。
catalog.properties.warehouse: your_warehouse
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
# テーブルにプライマリキー情報を追加します。
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# inventory トピックのすべてのデータを test_database.inventory テーブルに書き込みます。
route:
- source-table: inventory
sink-table: test_database.inventory
pipeline:
# (オプション) 処理例外を引き起こすダーティデータをログに記録します。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerApache Kafka の inventory トピックが、JSON フォーマットの複数のログテーブルのデータを格納していると仮定します。JSON ペイロードの databaseName フィールドと tableName フィールドが、データベース名とテーブル名を提供します。次のサンプルジョブは、これらのテーブルのデータを DLF の対応する結果テーブルに同期します。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
json.infer-schema.flatten-nested-columns.enable: true
# databaseName フィールドの値をデータベース名として、tableName フィールドの値をテーブル名として使用します。
json.decode.parser-table-id.fields: databaseName,tableName
# (オプション) 最初の 100 件の解析例外をスキップします。例外の数が 100 を超えると、ジョブは失敗します。
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
# メタストアのタイプ。値は rest に固定されます。
catalog.properties.metastore: rest
# トークンプロバイダー。値は dlf に固定されます。
catalog.properties.token.provider: dlf
# DLF Rest Catalog Server にアクセスするための URI。フォーマットは http://[region-id]-vpc.dlf.aliyuncs.com です。例: http://cn-hangzhou-vpc.dlf.aliyuncs.com。
catalog.properties.uri: dlf_uri
# DLF カタログの名前。
catalog.properties.warehouse: your_warehouse
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
# テーブルにプライマリキー情報を追加します。
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# ods.inventory、ods.customer、ods.user のデータを、それぞれ test_database.inventory、test_database.customer、test_database.user テーブルに書き込みます。
route:
- source-table: ods.inventory
sink-table: test_database.inventory
- source-table: ods.customer
sink-table: test_database.customer
- source-table: ods.user
sink-table: test_database.user
pipeline:
# (オプション) 処理例外を引き起こすダーティデータをログに記録します。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerJSON フォーマットの Kafka ソーステーブルのスキーマ解析および進化ポリシーの詳細については、「スキーマの解析と進化のポリシー」をご参照ください。
ユースケース
以下のセクションでは、一般的なユースケースのジョブ設定について説明します。より詳細な設定については、「データインジェスト Kafka コネクタ」をご参照ください。
フィールド名の競合の解決
Apache Kafka のメッセージは、キーと値で構成されます。key.format と value.format を設定することで、各部分のフォーマットを定義できます。最終的なスキーマは、両方の部分のすべてのフィールドを組み合わせたものになります。
キー部分と値部分の間でフィールド名が競合する場合は、key.fields-prefix と value.fields-prefix を使用してフィールド名にプレフィックスを追加します。
例えば、Kafka メッセージのキー部分に id と name フィールドが含まれ、値部分に id と price フィールドが含まれる場合、以下のジョブ設定では key_id、key_name、val_id、val_price というフィールドを持つスキーマが生成されます。
source:
type: kafka
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
key.format: json
# キーのフィールド名にプレフィックスを追加します。
key.fields-prefix: key_
# 値のフィールド名にプレフィックスを追加します。
value.fields-prefix: val_
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerメタデータの読み取り
metadata.list 設定を使用して、追加の Kafka メッセージメタデータを読み取り、渡すことができます。metadata.list に追加されたメタデータ列は、transform モジュールで直接使用できます。サポートされているメタデータのリストについては、「利用可能なメタデータ列」をご参照ください。
以下の設定では、partition と offset のメタデータをデータに追加します。その後、このメタデータを transform モジュールで使用して、パーティションが 1 より大きく、オフセットが 100 より大きいデータをフィルターします。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# メタデータ列を追加します。
metadata.list: partition,offset
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
transform:
- source-table: \.*.\.*
filter: '`partition` > 1 and `offset` > 100'
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerダーティデータの処理
ログデータには、不正なフォーマットのダーティデータが含まれている可能性があり、ジョブの失敗と再起動を繰り返す原因となることがあります。Flink CDC データインジェストは、解析エラーを無視し、解析に失敗したデータを収集することをサポートしています。詳細については、「ダーティデータ収集」をご参照ください。
以下のジョブは、解析に失敗したダーティデータをログファイルに書き込みます。100 件を超える解析エラーが発生した場合、ジョブは失敗します。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerダーティデータを収集する必要がなく、ダーティデータによるジョブの失敗を防ぎたい場合は、json.ignore-parse-errors、debezium-json.ignore-parse-errors、または canal-json.ignore-parse-errors を使用して解析エラーを直接無視できます。
TableID の指定
Debezium JSON および Canal JSON フォーマットは固定されており、特定のフィールドが TableID を格納するために使用されます。固定フォーマットを持たない一般的な JSON データの場合、デフォルトではトピック名が TableID として機能します。データから特定のフィールドを TableID として使用するには、json.decode.parser-table-id.fields を設定します。例えば、JSON データ {"col0":"a", "col1":"b", "col2":"c"} がある場合、異なる設定で以下の TableID が生成されます。
設定 | TableID |
col0 | a |
col0,col1 | a.b |
col0,col1,col2 | a.b.c |
以下のジョブ設定では、データから db 列と table 列を連結して TableID を作成します。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# データの db フィールドと table フィールドを TableID として使用します。
json.decode.parser-table-id.fields: db,table
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerデータ型の推論
ログデータのデータ型は、データを解析することによって推論されます。型推論の詳細については、「スキーマの解析と進化のポリシー」をご参照ください。以下のセクションでは、一般的なシナリオで型解析を制御するための設定調整方法について説明します。
(一般) すべてのフィールド型を文字列に設定
下流の処理やストレージが特定のフィールド型を必要としない場合は、json.infer-schema.primitive-as-string、debezium-json.infer-schema.primitive-as-string、または canal-json.infer-schema.primitive-as-string を有効にできます。これにより、型推論がスキップされ、すべてのフィールド型が String に設定されます。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# JSON データ内のすべてのフィールドの型を String に設定します。
json.infer-schema.primitive-as-string: true
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger(一般) 初期スキーマの指定
一部のシナリオでは、Kafka データを既存の下流テーブルに書き込む場合など、初期スキーマを指定する必要があります。scan.value.initial-schemas.ddls パラメーターを追加することで、初期スキーマを指定できます。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# データの db フィールドと table フィールドを TableID として使用します。
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);
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger上記の設定では、db1.t1 テーブルの id フィールドの初期型を BIGINT に、name フィールドを VARCHAR(10) に指定します。また、db1.t2 テーブルの id フィールドの初期型を BIGINT に指定します。
(JSON) 固定フィールド型の設定
JSON データの場合、フィールド型は JSON ノードタイプを解析して推論されます。推論された型が期待通りでない場合があります。これを解決するには、json.infer-schema.fixed-types 設定を使用して、特定のフィールドの型を指定します。
以下のジョブ設定では、id フィールドの型を BIGINT に、name フィールドの型を VARCHAR(10) に指定します。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# 特定のフィールドに固定の型を設定します。
json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
# バージョン 11.5 以前で必須です。
scan.max.pre.fetch.records: 0
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger(Canal JSON) 推論ソースの指定
Canal JSON データは、標準の JSON データよりも多くの情報を含んでいます。Canal JSON データに sqlType または mysqlType フィールドが存在する場合、この情報を使用してより正確なデータ型を解析できます。
sqlTypeからのスキーマ解析
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: canal-json
# sqlType フィールドからスキーマを推論します。
canal-json.infer-schema.strategy: SQL_TYPE
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggermysqlTypeからのスキーマ解析source: type: kafka name: Kafka source properties.bootstrap.servers: localhost:9092 topic: test_topic properties.group.id: test_group scan.startup.mode: earliest-offset value.format: canal-json # mysqlType フィールドからスキーマを推論します。 canal-json.infer-schema.strategy: MYSQL_TYPE # スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。 schema.inference.strategy: continuous # データ解析中のエラーを無視します。 ingestion.ignore-errors: true # 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。 ingestion.error-tolerance.max-count: 100 sink: type: paimon catalog.properties.metastore: rest catalog.properties.uri: dlf_uri catalog.properties.warehouse: your_warehouse catalog.properties.token.provider: dlf # (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。 table.properties.deletion-vectors.enabled: true pipeline: # ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。 dirty-data.collector: name: Logger Dirty Data Collector type: logger
静的スキーマの使用
トピック内のデータのスキーマが固定されている場合は、schema.inference.strategy を static に設定します。データインジェストジョブは、起動時に一度だけスキーマ推論を実行し、後続のデータのスキーマは解析しません。
source:
type: kafka
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# スキーマ推論戦略を static に設定します。スキーマはジョブの開始時に一度だけ推論されます。
schema.inference.strategy: static
# 各パーティションから 20 レコードを消費してスキーマを推論しようとします。
scan.max.pre.fetch.records: 20
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
# test_topic トピックのすべてのデータを test_database.test_topic テーブルに書き込みます。
route:
- source-table: test_topic
sink-table: test_database.test_topic
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerまた、scan.value.initial-schemas.ddls パラメーターを追加して初期スキーマを指定し、特定のテーブルのスキーマ推論をスキップすることもできます。次の例では、db1.t1 テーブルと db1.t2 テーブルの初期スキーマを指定します。
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: test_topic
properties.group.id: test_group
scan.startup.mode: earliest-offset
value.format: json
# データの db フィールドと table フィールドを TableID として使用します。
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);
# スキーマ推論戦略を static に設定します。スキーマはジョブの開始時に一度だけ推論されます。
schema.inference.strategy: static
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerKafka ログデータ解析の高速化
一般的な Kafka コネクタの高速化設定に加えて、Flink CDC データインジェストジョブは、ユースケースに基づいて解析を高速化するための特定の設定を提供します。
スキーマ解析には時間がかかることがあります。下流のアプリケーションが String 型のみを必要とする場合は、
json.infer-schema.primitive-as-string、debezium-json.infer-schema.primitive-as-string、またはcanal-json.infer-schema.primitive-as-stringを有効にできます。これにより、型推論がスキップされ、すべてのフィールド型が String に設定されるため、解析が高速化されます。Canal JSON データの場合、
canal-json.database.includeとcanal-json.table.includeを使用して、不要なテーブルのデータをフィルターできます。データスキーマが変更されない場合、またはスキーマ進化が必要ない場合は、
schema.inference.strategyをstaticに変更できます。これにより、ジョブの開始時に一度だけスキーマ推論が実行されます。
データレイクへの SLS ログのインジェスト
Log Service (SLS) は、ログデータのためのワンストップサービスです。Flink CDC データインジェストを使用すると、SLS からのログデータをリアルタイムで迅速にデータレイクにインジェストできます。自動的にスキーマ推論を実行し、スキーマ進化をサポートします。
source:
type: sls
endpoint: localhost
project: test_pj
logstore: test_log
accessId: access_id
accessKey: access_key
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerTableID:デフォルトでは、
projectとlogstoreが連結されて TableID が形成されます。上記のジョブでは、TableID はtest_pj.test_logになります。データ型:SLS コネクタは、デフォルトで各ログエントリのすべてのフィールドを String 型として扱います。
スキーマ進化:スキーマ進化は現在、新しい列の追加に限定されています。新しい列は、デフォルトでスキーマの末尾に追加されます。
ユースケース
以下のセクションでは、一般的なユースケースのジョブ設定について説明します。より詳細な設定については、「データインジェスト SLS コネクタ」をご参照ください。
メタデータの読み取り
metadata.list 設定を使用して、追加の SLS メタデータを読み取り、渡すことができます。metadata.list に追加されたメタデータ列は、transform モジュールで直接使用できます。サポートされているメタデータのリストについては、「metadata.list」をご参照ください。
metadata.list で追加された列は、自動的に出力データに含まれないことに注意してください。これらのメタデータ列を sink に書き込むには、transform モジュールの projection セクションで宣言する必要があります。
以下の設定では、__timestamp__ と __tag__ のメタデータをデータに追加し、それらを sink に書き込みます。また、これらのメタデータフィールドを transform モジュールで使用して、タイムスタンプが 1772181154 より大きく、タグが "test" のデータをフィルターします。
source:
type: sls
endpoint: localhost
project: test_pj
logstore: test_log
accessId: access_id
accessKey: access_key
# メタデータ列を追加します。
metadata.list: __timestamp__,__tag__
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
transform:
- source-table: \.*.\.*
projection: \*, __timestamp__ as timestamp_col, __tag__ as tag_col
filter: '`__timestamp__` > 1772181154 and `__tag__` = "test"'
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerダーティデータの処理
ログデータには、不正なフォーマットのダーティデータが含まれている可能性があり、ジョブの失敗と再起動を繰り返す原因となることがあります。Flink CDC データインジェストは、解析エラーを無視し、エラーの原因となったデータを収集することをサポートしています。詳細については、「ダーティデータ収集」をご参照ください。
以下のジョブは、解析に失敗したダーティデータをログファイルに書き込みます。100 件を超える解析エラーが発生した場合、ジョブは失敗します。
source:
type: sls
endpoint: localhost
project: test_pj
logstore: test_log
accessId: access_id
accessKey: access_key
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerTableID の指定
データから特定のフィールドを TableID として使用するには、decode.table-id.fields を設定します。例えば、ログデータ {"col0":"a", "col1":"b", "col2":"c"} がある場合、異なる設定で以下の TableID が生成されます。
設定 | TableID |
col0 | a |
col0,col1 | a.b |
col0,col1,col2 | a.b.c |
以下のジョブ設定では、データから db 列と table 列を連結して TableID を作成します。
source:
type: sls
endpoint: localhost
project: test_pj
logstore: test_log
accessId: access_id
accessKey: access_key
# データの db フィールドと table フィールドを TableID として使用します。
decode.table-id.fields: db,table
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: loggerフィールド型の指定
データインジェスト SLS コネクタは、デフォルトでデータ内のすべてのフィールドを String 型として扱います。特定のフィールドの型を指定する必要がある場合は、fixed-types 設定オプションを使用します。
以下のジョブでは、id フィールドの型を BIGINT に、name フィールドの型を VARCHAR(10) に指定します。
source:
type: sls
endpoint: localhost
project: test_pj
logstore: test_log
accessId: access_id
accessKey: access_key
# id フィールドの型を BIGINT に、name フィールドの型を VARCHAR(10) に指定します。
fixed-types: id BIGINT, name VARCHAR(10)
# スキーマ推論戦略を continuous に設定します。この戦略は、各メッセージのスキーマを検出し、スキーマの変更を同期します。
schema.inference.strategy: continuous
# データ解析中のエラーを無視します。
ingestion.ignore-errors: true
# 100 件のデータ解析エラーが発生した場合、ジョブは失敗します。
ingestion.error-tolerance.max-count: 100
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (オプション) 削除ベクターを有効にして、読み取りパフォーマンスを向上させます。
table.properties.deletion-vectors.enabled: true
pipeline:
# ダーティデータ収集を有効にします。ダーティデータはログファイルに書き込まれます。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger