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

Realtime Compute for Apache Flink:データインジェスト用 Hologres YAML コネクタ

最終更新日:Sep 10, 2026

このトピックでは、Hologres コネクタを使用して、YAML データインジェストのデプロイメントでデータを同期する方法について説明します。

背景

Hologres は、大規模なリアルタイムの書き込み、更新、分析をサポートするリアルタイムデータウェアハウスエンジンです。Hologres は PostgreSQL プロトコルと互換性があり、標準 SQL をサポートしています。ペタバイト規模のデータに対するオンライン分析処理 (OLAP) とアドホッククエリをサポートし、高い同時実行性と低レイテンシーでデータを提供します。Hologres は MaxCompute、Realtime Compute for Apache Flink、DataWorks と統合し、完全なオンラインおよびオフラインのデータウェアハウスソリューションを提供します。次の表に、Hologres YAML コネクタの機能を示します。

項目

説明

テーブルタイプ

シンク

実行モード

ストリーミングおよびバッチモード

データフォーマット

N/A

メトリック

  • numRecordsOut

  • numRecordsOutPerSecond

説明

詳細については、「メトリック」をご参照ください。

API タイプ

YAML

シンクテーブルでの更新または削除

サポート

機能

機能

説明

データベース全体の同期

データベース全体または複数のテーブルから、対応するシンクテーブルへ、フルデータおよび増分データをリアルタイムで同期します。

スキーマ変更の同期

ソーステーブルから対応するシンクテーブルへ、スキーマの変更 (列の追加、削除、名前の変更など) をリアルタイムで同期します。

シャーディングされたデータベースとテーブルの同期

正規表現を使用して、複数のシャーディングされたデータベース全体で名前でソーステーブルを照合します。これらのテーブルからのデータは、マージされ、名前が対応するダウンストリームのシンクテーブルに同期されます。

パーティションテーブルへの書き込み

アップストリームテーブルから Hologres パーティションテーブルへデータを書き込みます。

データ型マッピング

複数の戦略を使用して、アップストリームのデータ型をより大きい Hologres データ型にマッピングします。

構文

sink:
  type: hologres
  name: Hologres シンク
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}

パラメータ

パラメーター

説明

タイプ

必須

デフォルト

備考

type

シンクのタイプ。

String

はい

なし

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

name

シンクの名前。

String

いいえ

なし

N/A。

dbname

データベース名。

String

はい

なし

N/A。

username

データベースアクセス用のユーザー名。Alibaba Cloud アカウントの AccessKey ID を使用します。

String

はい

なし

詳細については、「AccessKey ID と AccessKey Secret の表示方法」をご参照ください。

重要

AccessKey が公開されるのを防ぐため、変数を使用してその値を指定します。詳細については、「プロジェクト変数」をご参照ください。

password

データベースアクセス用のパスワード。Alibaba Cloud アカウントの AccessKey Secret を使用します。

String

はい

なし

endpoint

Hologres サービスのエンドポイント。

String

はい

なし

詳細については、「アクセスエンドポイント」をご参照ください。

jdbcRetryCount

接続が失敗した場合の書き込みおよびクエリ操作のリトライ回数。

Integer

いいえ

10

N/A。

jdbcRetrySleepInitMs

各リトライ試行の固定待機時間。

Long

いいえ

1000

単位: ミリ秒。 再試行の実際の待機時間は、次の計算式で算出されます。 jdbcRetrySleepInitMs + retry * jdbcRetrySleepStepMs。

jdbcRetrySleepStepMs

各リトライ試行の増分待機時間。

Long

いいえ

5000

単位:ミリ秒。リトライの実際の待機時間は、次の計算式を用いて算出されます: jdbcRetrySleepInitMs + retry * jdbcRetrySleepStepMs。

jdbcConnectionMaxIdleMs

JDBC 接続の最大アイドル時間。

Long

いいえ

60000

単位:ミリ秒。接続がこの期間を超えてアイドル状態になると、切断され、解放されます。

jdbcMetaCacheTTL

ローカルにキャッシュされた TableSchema 情報の有効期限。

Long

いいえ

60000

単位:ミリ秒。

jdbcMetaAutoRefreshFactor

キャッシュ更新トリガーを決定する係数。キャッシュの残りの有効期間がトリガー時間より短い場合、システムは自動的にキャッシュを更新します。

Integer

いいえ

4

キャッシュの残り有効期間は、次の式で計算されます。 キャッシュの残り有効期間 = キャッシュの有効期限 - キャッシュのアクティブ時間。自動更新後、キャッシュのアクティブ時間は 0 にリセットされます。

トリガー時刻は、式 jdbcMetaCacheTTL / jdbcMetaAutoRefreshFactor を使用して計算されます。

mutatetype

データ書き込みモード。

String

いいえ

INSERT_OR_UPDATE

Hologres 物理テーブルにプライマリキーが設定されている場合、Hologres シンクはプライマリキーに基づいて exactly-once セマンティクスを保証します。プライマリキーが重複するデータが到着した場合、 mutatetype パラメーターはシンクテーブルの更新方法を決定します。 mutatetype パラメーターは、次の値をサポートします。

  • INSERT_OR_IGNORE:最初のデータを保持し、後続のすべてのデータを無視します。

  • INSERT_OR_REPLACE:新しいデータが既存の行全体を置き換えます。

  • INSERT_OR_UPDATE:既存の行の列のサブセットを更新します。例えば、テーブルに a、b、c、d の 4 つの列があり、a がプライマリキー (PK) であるとします。列 a と b のデータのみを Hologres に書き込み、同じ PK を持つ行が既に存在する場合、システムは列 b のみを更新します。列 c と d は変更されません。

createparttable

パーティションテーブルに書き込む際に、存在しないパーティションを自動的に作成するかどうかを指定します。

Boolean

いいえ

false

N/A。

sink.delete-strategy

取り消しメッセージの処理方法を指定します。

String

いいえ

なし

有効な値:

  • IGNORE_DELETE:Update Before および Delete メッセージを無視します。この値は、データの挿入または更新のみが必要で、削除は不要なシナリオに適しています。

  • DELETE_ROW_ON_PK:Flink フレームワークは、プライマリキーに基づいて削除操作を適用します。更新操作の場合、データの正確性を保証するために、まず古いデータを削除し、次に新しいデータを挿入します。

jdbcWriteBatchSize

JDBC モードにおいて、バッチ書き込みの前に Hologres シンクでバッファリングするレコードの最大数。

Integer

いいえ

256

単位:行。

説明

jdbcWriteBatchSize、jdbcWriteBatchByteSize、および jdbcWriteFlushInterval パラメーターは OR 関係にあります。3 つすべてのパラメーターを設定した場合、いずれかの条件が満たされるとすぐに結果データが書き込まれます。

jdbcWriteBatchByteSize

JDBC モードにおいて、このパラメーターは、バッチを宛先に書き込む前に Hologres シンクがバッファリングするデータの最大サイズ (バイト単位) を指定します。

Long

いいえ

2,097,152 バイト (2 MB)

説明

jdbcWriteBatchSize、jdbcWriteBatchByteSize、および jdbcWriteFlushInterval パラメーターは OR 関係にあります。3 つすべてのパラメーターを設定した場合、いずれかの条件が満たされるとすぐに結果データが書き込まれます。

jdbcWriteFlushInterval

JDBC モードにおいて、このパラメーターは、バッファリングされたデータを Hologres に書き込む前に Hologres シンクが待機する最大時間を指定します。

Long

いいえ

10000

単位:ミリ秒。

説明

jdbcWriteBatchSize、jdbcWriteBatchByteSize、および jdbcWriteFlushInterval パラメーターは OR 関係にあります。3 つすべてのパラメーターを設定した場合、いずれかの条件が満たされるとすぐに結果データが書き込まれます。

ignoreNullWhenUpdate

mutatetype が INSERT_OR_UPDATE に設定されている場合に、入力データ内の null 値を無視するかどうかを指定します。

Boolean

いいえ

false

有効な値:

  • false (デフォルト):Hologres 結果テーブルに null 値を書き込みます。

  • true:受信データ内の null 値を無視します。

jdbcEnableDefaultForNotNullColumn

定義されたデフォルト値がない NOT NULL 列に null が書き込まれた場合に、デフォルト値を挿入するかどうかを指定します。

Boolean

いいえ

true

有効な値:

  • true (デフォルト):コネクタがデフォルト値を挿入できるようにします。ルールは次のとおりです:

    • String 型の列には、空の文字列 ("") が書き込まれます。

    • Number 型の列には、0 が書き込まれます。

    • Date、timestamp、または timestamptz 型の列には、 1970-01-01 00:00:00 が書き込まれます。

  • false:デフォルト値を挿入しません。NOT NULL 列に null 値が書き込まれると、例外がスローされます。

remove-u0000-in-text.enabled

書き込み前に文字列から NULL 文字 (\u0000) を削除するかどうかを指定します。

Boolean

いいえ

false

有効な値:

  • false (デフォルト):コネクタはデータを処理しません。ただし、ダーティデータが検出された場合、書き込み操作で次の例外がスローされることがあります: ERROR: invalid byte sequence for encoding "UTF8": 0x00

    この場合、ソーステーブル内のダーティデータを処理するか、SQL ステートメントでデータ処理ロジックを定義する必要があります。

  • true:コネクタは書き込み例外を防ぐために、文字列から \u0000 文字を削除します。

deduplication.enabled

jdbc および jdbc_fixed モードでの書き込み前に、各バッチ内で重複排除を実行するかどうかを指定します。

Boolean

いいえ

true

有効な値:

  • true (デフォルト):バッチ内での重複排除を有効にします。複数のレコードが同じプライマリキーを持つ場合、最後のレコードのみが保持されます。例えば、最初のフィールドがプライマリキーである 2 つのフィールドを持つデータを考えます:

    • レコード INSERT (1,'a') と INSERT (1,'b') が順次到着します。重複排除後、後から到着したレコード (1,'b') のみが Hologres 結果テーブルに書き込まれます。

    • Hologres 結果テーブルに既にレコード (1,'a') が含まれているとします。DELETE (1,'a') と INSERT (1,'b') のレコードが順次到着した場合、最後に到着したレコード (1,'b') のみが Hologres に書き込まれます。これは、削除とそれに続く挿入ではなく、直接の更新として扱われます。

  • false:バッチ処理中に重複排除は実行されません。新しいレコードが現在のバッチに既にあるレコードと同じプライマリキーを持つ場合、現在のバッチが先に書き込まれます。書き込みが完了した後、新しいレコードが処理されます。

sink.type-normalize-strategy

データ型マッピング戦略。

String

いいえ

STANDARD

Hologres シンクがアップストリームのデータ型を Hologres の型に変換するために使用する戦略。

  • STANDARD:Flink CDC の型を標準の PostgreSQL (PG) 型に変換します。

  • BROADEN:Flink CDC の型をより広範囲の Hologres 型に変換します。

  • ONLY_BIGINT_OR_TEXT:すべての Flink CDC 型を Hologres の BIGINT または TEXT のいずれかに変換します。

sink.insert.legacy-put-handler

従来の Put Handler を使用して Hologres にデータを書き込むかどうかを指定します。

Boolean

いいえ

false

有効な値:

  • false (デフォルト):新しい Put Handler を使用してデータを書き込みます。書き込み操作の SQL 形式は insert into xxx(c0,c1,...) select unnest(?),unnest(?),... on conflict です。

  • true: レガシー Put ハンドラーを使用してデータを書き込みます。書き込みの SQL 形式は insert into xxx(c0,c1,...) values (?,?,...),... on conflict;  です。

table_property.*

Hologres の物理テーブルプロパティ。

String

いいえ

なし

Hologres テーブルを作成する際、WITH 句で物理テーブルプロパティを設定できます。適切なテーブルプロパティは、システムがデータを効率的に整理し、クエリを実行するのに役立ちます。

警告

table_property.distribution_key パラメーターのデフォルトはプライマリキーの値です。データ書き込みの正確性に影響を与える可能性があるため、その影響を完全に理解していない限り、この設定を変更しないでください。

connection.ssl.mode

Secure Sockets Layer (SSL) 転送時の暗号化を有効にするかどうか、および使用するモードを指定します。

String

いいえ

disable

  • disable (デフォルト):転送時の暗号化を無効にします。

  • require:SSL を有効にしてデータリンクを暗号化します。

  • verify-ca:SSL を有効にし、データリンクを暗号化し、CA 証明書を使用して Hologres サーバーの信頼性を検証します。

  • verify-full:SSL を有効にし、データリンクを暗号化し、CA 証明書を使用して Hologres サーバーの信頼性を検証し、証明書内の CN または DNS 名が設定された Hologres エンドポイントと一致することを確認します。

説明
  • Hologres V2.1 以降では、verify-ca および verify-full モードがサポートされています。詳細については、「転送時の暗号化」をご参照ください。

  • このパラメーターを verify-ca または verify-full に設定する場合は、 connection.ssl.root-cert.location パラメーターも設定する必要があります。

connection.ssl.root-cert.location

転送時の暗号化モードで証明書が必要な場合に、証明書ファイルへのパスを指定します。

String

いいえ

なし

connection.ssl.mode が verify-ca または verify-full に設定されている場合、CA 証明書へのパスも設定する必要があります。 Realtime Compute コンソールのファイル管理機能を使用して、証明書をプラットフォームにアップロードできます。 証明書がアップロードされると、/flink/usrlib ディレクトリに保存されます。 たとえば、CA 証明書ファイルの名前が certificate.crt の場合、パラメーター値は '/flink/usrlib/certificate.crt' にする必要があります。

説明

CA 証明書を取得するには、「転送時の暗号化 - CA 証明書のダウンロード」をご参照ください。

connection.akv4.enabled

AKV4 モードを有効にして Hologres サーバーに接続するかどうかを指定します。

Boolean

いいえ

false

N/A。

connection.akv4.region

AKV4 モードが有効な場合に、サーバーが配置されているリージョンを指定します。

String

いいえ

なし

例えば、 cn-shanghai です。

既存のカタログの再利用

VVR 11.5 以降では、データ管理ページで作成された組み込みの Hologres カタログを Flink CDC データインジェスト ジョブで直接参照できます。これにより、接続プロパティを手動で指定する手間が省けます。

sink:
  type: hologres
  using.built-in-catalog: my_holo_catalog

データインジェスト ジョブは、次の Hologres カタログパラメーターを自動的に再利用できます:

  • エンドポイント

  • ユーザー名

  • パスワード

  • データベース名

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

データ型マッピング

パラメーター sink.type-normalize-strategy を使用して、アップストリームデータを Hologres 型に変換する戦略を設定できます。

説明
  • YAML ジョブを初めて起動するときに sink.type-normalize-strategy を有効にしてください。ジョブの開始後に有効にする場合は、ダウンストリームテーブルを削除し、ステートなしでジョブを再起動しないと、設定が有効になりません。

  • 現在、配列型は INTEGER、BIGINT、FLOAT、DOUBLE、BOOLEAN、CHAR、VARCHAR のみをサポートします。

  • Hologres は numeric 型をプライマリキーとしてサポートしていません。プライマリキーの型が numeric にマッピングされる場合、システムはそれを varchar 型に変換します。

STANDARD

sink.type-normalize-strategy が STANDARD に設定されている場合、型は次のようにマッピングされます:

Flink CDC 型

Hologres 型

CHAR

bpchar

STRING

text

VARCHAR

text (長さ > 10485760 の場合)

varchar (長さ <= 10485760 の場合)

BOOLEAN

bool

BINARY

bytea

VARBINARY

DECIMAL

numeric

TINYINT

int2

SMALLINT

INTEGER

int4

BIGINT

int8

FLOAT

float4

DOUBLE

float8

DATE

date

TIME_WITHOUT_TIME_ZONE

time

TIMESTAMP_WITHOUT_TIME_ZONE

timestamp

TIMESTAMP_WITH_LOCAL_TIME_ZONE

timestamptz

ARRAY

対応する要素型の配列

MAP

非対応

ROW

非対応

BROADEN

sink.type-normalize-strategy が BROADEN に設定されている場合、Flink CDC 型はより広範囲の Hologres 型に変換されます。型は次のようにマッピングされます:

Flink CDC 型

Hologres 型

CHAR

text

STRING

VARCHAR

BOOLEAN

bool

BINARY

bytea

VARBINARY

DECIMAL

numeric

TINYINT

int8

SMALLINT

INTEGER

BIGINT

FLOAT

float8

DOUBLE

DATE

date

TIME_WITHOUT_TIME_ZONE

time

TIMESTAMP_WITHOUT_TIME_ZONE

timestamp

TIMESTAMP_WITH_LOCAL_TIME_ZONE

timestamptz

ARRAY

対応する要素型の配列

MAP

非対応

ROW

非対応

ONLY_BIGINT_OR_TEXT

sink.type-normalize-strategy が ONLY_BIGINT_OR_TEXT に設定されている場合、すべての Flink CDC 型は Hologres の BIGINT または text 型に変換されます。データ型マッピングは次のとおりです:

Flink CDC 型

Hologres 型

TINYINT

int8

SMALLINT

INTEGER

BIGINT

BOOLEAN

text

BINARY

VARBINARY

DECIMAL

FLOAT

DOUBLE

DATE

TIME_WITHOUT_TIME_ZONE

TIMESTAMP_WITHOUT_TIME_ZONE

TIMESTAMP_WITH_LOCAL_TIME_ZONE

ARRAY

対応する要素型の配列

MAP

非対応

ROW

非対応

パーティションテーブルへの書き込み

Hologres シンクと変換を組み合わせ、アップストリームデータを Hologres パーティションテーブルに書き込むことができます。

  • パーティションキー は プライマリキー の一部である必要があります。アップストリームデータの非プライマリキー列を パーティションキー として使用すると、アップストリームテーブルとダウンストリームテーブルのプライマリキーに不整合が生じ、データ同期中にデータの不一致が発生する可能性があります。

  • Hologres では、TEXT、VARCHAR、INT のデータ型の列を パーティションキー として使用できます。バージョン 1.3.22 以降では、DATE データ型の列もサポートされています。

  • 子パーティションテーブルを自動的に作成するには、createparttable パラメーターを true に設定します。設定しない場合は、手動で作成する必要があります。

例については、「Writing data to a partitioned table」をご参照ください。

テーブルスキーマ同期

CDC YAML パイプラインは、テーブルスキーマの変更を処理するためにさまざまな戦略を使用します。これらは、パイプラインレベルのパラメーター schema.change.behavior で設定できます。schema.change.behavior の有効な値は、IGNORE、LENIENT、TRY_EVOLVE、EVOLVE、および EXCEPTION です。Hologres シンクは現在、TRY_EVOLVE 戦略をサポートしていません。LENIENT および EVOLVE 戦略は、テーブルスキーマの変更を伴います。以降のセクションでは、これら 2 つのモードがさまざまなスキーマ変更イベントをどのように処理するかを説明します。

LENIENT (デフォルト)

LENIENT モードでは、スキーマの変更は次のように処理されます。

  • NULL 許容列の追加:対応する列がシンクテーブルの末尾に自動的に追加され、そのデータが同期されます。

  • NULL 許容列の削除:列はシンクテーブルから削除されません。代わりに、その列には自動的に NULL 値が設定されます。

  • NOT NULL 列の追加:対応する NULL 許容列がシンクテーブルの末尾に自動的に追加され、そのデータが同期されます。既存の行では、この新しい列に自動的に NULL が設定されます。

  • 列の名前変更:この操作は、列の削除と新しい列の追加として扱われます。指定された名前の新しい列がシンクテーブルの末尾に追加され、元の名前の列には自動的に NULL 値が設定されます。例えば、col_a が col_b に名前変更された場合、col_b 列がシンクテーブルの末尾に追加され、col_a 列には自動的に NULL 値が設定されます。

  • 列のデータ型の変更:サポートされていません。Hologres は列のデータ型の変更をサポートしていないため、sink.type-normalize-strategy パラメーターを使用する必要があります。

  • 以下のスキーマ変更はサポートされていません。

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

    • NOT NULL 列の削除。

    • 列を NOT NULL から NULLABLE に変更。

EVOLVE

EVOLVE モードでは、スキーマの変更は次のように処理されます。

  • NULL 許容列の追加:サポートされています。

  • NULL 許容列の削除:サポートされていません。

  • NOT NULL 列の追加:新しい NULL 許容列がシンクテーブルに追加されます。

  • 列の名前変更:サポートされています。元の列はシンクテーブルで名前が変更されます。

  • 列のデータ型の変更:サポートされていません。Hologres は列のデータ型の変更をサポートしていないため、sink.type-normalize-strategy パラメーターを使用する必要があります。

  • 以下のスキーマ変更はサポートされていません。

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

    • NOT NULL 列の削除。

    • 列を NOT NULL から NULLABLE に変更。

警告

EVOLVE モードでは、シンクテーブルを削除せずにステートレス再起動を実行すると、アップストリームデータとシンクテーブル間のスキーマの不整合によりパイプラインが失敗する可能性があります。その場合、手動でシンクテーブルのスキーマを調整する必要があります。

例については、「EVOLVE モードの有効化」をご参照ください。

コード例

型拡張

sink.type-normalize-strategy パラメーターを使用して型拡張を設定します。

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

パーティションテーブルへの書き込み

create_time タイムスタンプフィールドを日付型に変換し、Hologres テーブルのパーティションキーとして使用します。

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # パーティションテーブルが存在しない場合は自動的に作成します。
  createparttable: true
 
transform:
  - source-table: test_db.test_source_table
    projection: \*, DATE_FORMAT(CAST(create_time AS TIMESTAMP), 'yyyy-MM-dd') as partition_key
    primary-keys: id, create_time, partition_key
    partition-keys: partition_key
    description: パーティションキーの追加 

pipeline:
  name: MySQL to Hologres Pipeline

EVOLVE モードの有効化

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # パーティションテーブルが存在しない場合は自動的に作成します。
  createparttable: true

pipeline:
  name: MySQL to Hologres Pipeline
  schema.change.behavior: evolve

単一テーブル同期

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

完全なデータベース同期

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

シャードテーブルのマージ

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.user\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。  
  sink.type-normalize-strategy: BROADEN
  
route:
  # MySQL の test_db データベース内のすべてのシャードテーブルを、test_db.user という名前の単一の Hologres テーブルにマージします。
  - source-table: test_db.user\.*
    sink-table: test_db.user

pipeline:
  name: MySQL to Hologres Pipeline

指定スキーマへの同期

Hologres では、スキーマは MySQL のデータベースに対応します。同期先テーブルのスキーマを指定できます。

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN
  
route:
  # 元のテーブル名を維持したまま、MySQL の test_db データベースからすべてのテーブルを Hologres の test_db2 スキーマに同期します。
  - source-table: test_db.\.*
    sink-table: test_db2.<>
    replace-symbol: <>

pipeline:
  name: MySQL to Hologres Pipeline

再起動なしでの新しいテーブルの同期

ジョブの実行中に新しく追加されたテーブルをリアルタイムで同期するには、scan.binlog.newly-added-table.enabled: true を設定します。

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499
  # ジョブの実行中に作成された新しいテーブルを自動的にキャプチャします。  
  scan.binlog.newly-added-table.enabled: true

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

再起動時に既存のテーブルを追加

同期に既存のテーブルを含めるには、scan.newly-added-table.enabled を true に設定してジョブを再起動します。

警告

以前に scan.binlog.newly-added-table.enabled で実行されたジョブで scan.newly-added-table.enabled を使用しないでください。この組み合わせは、再起動時にデータの重複を引き起こします。

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499
  scan.startup.mode: initial
  # 再起動時に、ジョブは `tables` パラメーターに一致する新しいテーブルをスキャンし、スナップショットを実行します。
  # 注:このパラメーターは scan.startup.mode: initial と一緒に使用する必要があります。
  scan.newly-added-table.enabled: true

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

テーブルの除外

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  # この正規表現に一致するテーブルを除外します。
  tables.exclude: test_db.table1
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # CDC データ型をより広範な Hologres 型にマッピングします。
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

関連ドキュメント