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

Realtime Compute for Apache Flink:データレイクへのリアルタイムログインジェスト

最終更新日:Apr 23, 2026

このトピックでは、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: logger

Apache 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: logger

JSON フォーマットの 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: logger
  • mysqlType からのスキーマ解析

    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: logger

Kafka ログデータ解析の高速化

一般的な Kafka コネクタの高速化設定に加えて、Flink CDC データインジェストジョブは、ユースケースに基づいて解析を高速化するための特定の設定を提供します。

  1. スキーマ解析には時間がかかることがあります。下流のアプリケーションが String 型のみを必要とする場合は、json.infer-schema.primitive-as-string、debezium-json.infer-schema.primitive-as-string、または canal-json.infer-schema.primitive-as-string を有効にできます。これにより、型推論がスキップされ、すべてのフィールド型が String に設定されるため、解析が高速化されます。

  2. Canal JSON データの場合、canal-json.database.include と canal-json.table.include を使用して、不要なテーブルのデータをフィルターできます。

  3. データスキーマが変更されない場合、またはスキーマ進化が必要ない場合は、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: logger
  • TableID:デフォルトでは、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: logger

TableID の指定

データから特定のフィールドを 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