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

Realtime Compute for Apache Flink:Hologres SQL コネクタ

最終更新日:Sep 08, 2026

Hologres コネクタを使用して、ストリームモードとバッチモードで Hologres テーブルの読み書きを行います。CDC、binlog 消費、部分更新をサポートしています。

概要

Hologres は、標準 SQL (PostgreSQL 互換)、ペタバイト規模の OLAP、アドホック分析、高同時実行・低レイテンシーのデータサービングをサポートする、統合されたリアルタイムデータウェアハウスです。MaxCompute、Realtime Compute for Apache Flink、DataWorks と統合されています。次の表に、コネクタの機能を示します。

カテゴリ

説明

サポートされる型

ソーステーブル、ディメンションテーブル、結果テーブル

実行モード

ストリームモードとバッチモード

データフォーマット

サポートされていません

特定の監視メトリクス

監視メトリクス

  • ソーステーブル:

    • numRecordsIn

    • numRecordsInPerSecond

  • 結果テーブル:

    • numRecordsOut

    • numRecordsOutPerSecond

    • currentSendTime

    説明

    メトリクスの詳細: 監視メトリクス。

API タイプ

DataStream と SQL

結果テーブルでのデータ更新または削除をサポート

はい

特徴

特徴

説明

Hologres データのリアルタイム消費

CDC モードと非 CDC モードの両方で、binlog の有無にかかわらず Hologres データの読み取りをサポートします。

完全消費と増分消費の統合

完全消費、増分消費、および完全消費と増分消費の統合をサポートします。

プライマリキーの競合処理

新しいデータを無視したり、行全体を置換したり、特定のフィールドのみを更新したりできます。

複数ストリームのマージと部分更新

行全体ではなく、変更された列のみを更新します。

パーティションテーブルからの binlog 消費 (ベータ版)

物理パーティションテーブルと論理パーティションテーブルの両方からの binlog 消費をサポートします。物理パーティションテーブルの場合、単一のジョブで新しく追加されたパーティションを含むすべてのパーティションを監視できます。

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

親テーブルへの書き込みをサポートし、対応する子パーティションを自動的に作成します。

単一テーブルまたはデータベース全体のリアルタイム同期

単一テーブルまたはデータベース全体のリアルタイム同期をサポートし、以下の主要な特徴があります:

  • ソーステーブルのスキーマ進化の自動検出:ソーステーブルのスキーマが変更されると、Hologres はこれらの変更をリアルタイムで結果テーブルに同期します。

  • スキーマ変更の自動処理:新しいデータが取り込まれると、Flink はデータを書き込む前にまず結果テーブルのスキーマを変更します。

詳細については、「リアルタイムデータベース同期のクイックスタート」をご参照ください。

制限事項と推奨事項

制限事項

  • 外部テーブルはサポート対象外:Hologres コネクタは、MaxCompute 外部テーブルなどの Hologres 外部テーブルをサポートしていません。

  • 時間型の制約:Hologres コネクタは現在、TIMESTAMP データのリアルタイム消費をサポートしていません。テーブルを作成する際は、TIMESTAMPTZ 型のみを使用してください。

  • ソーステーブルのスキャンモード (Ververica Runtime (VVR) 8 以前):デフォルトでは、コネクタはバッチモードでデータを読み取り、テーブル全体を一度だけスキャンします。そのため、新しく追加されたデータは消費されません。

  • ウォーターマークの制約 (Ververica Runtime (VVR) 8 以前):CDC モードはウォーターマーク定義をサポートしていません。ウィンドウ集約を実行するには、代わりに非ウィンドウ集約ソリューションを使用してください。

  • COPY モードでの時間関数の動作:COPY_STREAM または COPY_BULK_LOAD 書き込みモードを使用する場合、デフォルト値として CURRENT_TIMESTAMP または NOW() を使用するテーブル列の値は、接続開始時刻に固定され、各行で更新されません。hg_binlog_timestamp_us binlog メタデータ列は、実際のインジェスト時間を提供します。

推奨事項

  • ストレージ形式の選択:

    • ディメンションテーブルのポイントルックアップの場合:行指向ストレージを使用します。プライマリキーとクラスタリングキーを設定する必要があります。

    • ディメンションテーブルでの 1 対多クエリの場合:列指向ストレージを使用します。最適なパフォーマンスを得るには、分散キーとセグメントキーを適切に設定してください。

    • 頻繁な更新と分析クエリが必要なテーブルの場合:テーブルがリアルタイムの binlog 消費と OLAP 分析の両方をサポートする必要がある場合、行列ハイブリッドストレージの使用を強く推奨します。

    重要

    Hologres でテーブルを作成する際にストレージ形式を指定しない場合、デフォルトで列指向ストレージになります。テーブル作成後にストレージ形式を変更することはできません。Hologres でのテーブル作成 | テーブルストレージ形式:行指向、列指向、行列ハイブリッド。

  • ジョブの並列度の設定:Flink ジョブの並列度を Hologres テーブルのシャード数に合わせて設定します。

    -- HoloWeb で、次のステートメントを実行してテーブルのシャード数を確認します。<tablename> をご利用のテーブル名に置き換えてください。
    select tg.property_value from hologres.hg_table_properties tb join hologres.hg_table_group_properties tg on tb.property_value = tg.tablegroup_name where tb.property_key = 'table_group' and tg.property_key = 'shard_count' and table_name = '<tablename>';
  • バージョンと機能:既知の問題、機能の更新、バージョンの互換性情報については、Hologres コネクタのリリースノートを定期的に確認してください。

注意事項

  • Hologres と VVR の消費モードの互換性と制限事項

    ソーステーブル

    • VVR 8 以前の場合、sdkMode パラメーターを使用して消費モードを選択します。

    • VVR 11 以降の場合、source.binlog.read-mode パラメーターを使用して消費モードを選択します。

    VVR バージョン

    Hologres バージョン

    デフォルト/推奨値

    実際の消費モード

    注意

    ≥ 6.0.7

    < 2.0

    カスタム

    holohub (デフォルト)

    値を jdbc に設定することを推奨します。

    6.0.7 から 8.0.4

    ≥ 2.0

    jdbc (自動切り替え、設定不要)

    jdbc (強制)

    holohub サービスは Hologres 2.0 以降で非推奨になりました。システムは自動的に jdbc モードに切り替わりますが、これにより権限の問題が発生する可能性があります。権限設定については、「権限の問題」をご参照ください。

    ≥ 8.0.5

    ≥ 2.1

    jdbc (自動切り替え、設定不要)

    jdbc (強制)

    権限の問題はありません。Hologres 2.1.27 以降では、モードは jdbc_fixed に切り替わります。

    ≥ 11.1

    すべてのバージョン

    AUTO (デフォルト)

    Hologres のバージョンに基づいて自動的に選択

    • Hologres 2.1.27 以降では、jdbc モードが選択され、デフォルトで軽量接続が有効になります (connection.fixed.enabled パラメーターが true に設定されます)。

    • Hologres バージョン 2.1.0 から 2.1.26 では、jdbc モードが選択されます。

    • Hologres 2.0 以前では、holohub モードが選択されます。

    重要

    VVR 11.1 以降では、コネクタはデフォルトで binlog データを消費します。binlog が有効になっていることを確認してください。そうでない場合、エラーが発生する可能性があります。

    権限の問題

    スーパーユーザーでない場合、JDBC モードで binlog を消費するには必要な権限を付与する必要があります。

    user_name は、ご利用の Alibaba Cloud アカウントまたは RAM ユーザーの ID です (アカウントの概要)。

    -- 標準 PostgreSQL 権限モデルでは、ユーザーに CREATE 権限を付与し、インスタンス上のレプリケーションロールをユーザーに付与します。
    GRANT CREATE ON DATABASE <db_name> TO <user_name>;
    alter role <user_name> replication;
    
    -- データベースが簡易権限モデル (SLMP) を使用している場合、GRANT ステートメントを実行できません。spm_grant を使用して、データベースの Admin 権限をユーザーに付与します。HoloWeb コンソールで権限を付与することもできます。
    call spm_grant('<db_name>_admin', '<user_name>');
    alter role <user_name> replication;

    シンクテーブル

    • VVR 8 以前の場合、sdkMode パラメーターを使用してデータ書き込みモードを選択します。

    • VVR 11 以降の場合、sink.write-mode パラメーターを使用してデータ書き込みモードを選択します。

    VVR バージョン

    Hologres バージョン

    影響を受ける RPC モード

    実際の書き込みモード

    デフォルト/推奨値

    注意

    6.0.4 から 8.0.2

    < 2.0

    いいえ

    rpc

    カスタム

    N/A

    6.0.4 から 8.0.2

    ≥ 2.0

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    重複排除を防ぐには、'jdbcWriteBatchSize'='1' を設定します。

    ≥ 8.0.3

    すべてのバージョン

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    モードを rpc として設定した場合、システムは自動的にモードを jdbc_fixed に切り替え、重複排除を防ぐために 'jdbcWriteBatchSize'='1' を設定します。

    ≥ 8.0.5

    すべてのバージョン

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    モードを rpc として設定した場合、システムは自動的にモードを jdbc_fixed に切り替え、重複排除を防ぐために 'deduplication.enabled'='false' を設定します。

    重要
    • rpc サービスは Hologres 2.0 以降で非推奨になりました。このパラメーターを rpc に設定すると、Flink は自動的に値を jdbc_fixed に切り替えます。パラメーターを別の値に設定した場合、Flink は指定された値を使用します。

    • rpc モードは VVR 11.1 以降から削除されました。接続には jdbc モードを使用することを推奨します。

    • 高同時実行シナリオでの書き込み操作には、jdbc_copy または COPY_STREAM モードを使用することを推奨します。

    ディメンションテーブル

    VVR バージョン

    Hologres バージョン

    影響を受ける RPC モード

    実際の消費モード

    デフォルト/推奨値

    注意

    6.0.4 から 8.0.2

    < 2.0

    いいえ

    rpc

    カスタム

    N/A

    6.0.4 から 8.0.2

    ≥ 2.0

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    rpc サービスは Hologres インスタンスのバージョン 2.0 以降で非推奨になりました。このパラメーターを rpc に設定すると、Flink は自動的に値を jdbc_fixed に切り替えます。ただし、パラメーターを別の値に設定した場合、Flink は指定した値を使用します。

    ≥ 8.0.3

    すべてのバージョン

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    ≥ 8.0.5

    すべてのバージョン

    はい

    jdbc_fixed (自動切り替え)

    カスタム

    重要

    rpc モードは VVR 11.1 以降から削除されました。デフォルトでは、接続に jdbc モードが使用されます。connection.fixed.enabled パラメーターを設定することで、軽量接続モードを有効にすることもできます。

  • JDBC モードで binlog ソーステーブルから JSONB データを読み取るには、データベースレベルで GUC パラメーターを有効にする必要があります。

    -- データベースレベルで GUC パラメーターを有効にします。このコマンドはスーパーユーザーのみが実行できます。各データベースに対してこのコマンドを一度だけ実行する必要があります。
    alter database <db_name> set hg_experimental_enable_binlog_jsonb = on;
  • UPDATE 操作は、2 つの連続した binlog レコードを生成します:古いデータのための update_before レコードと、それに続く新しいデータのための update_after レコードです。

  • binlog ソーステーブルで TRUNCATE やその他のテーブル再構築操作を実行しないでください。詳細については、「よくある質問」をご参照ください。

  • エラーを避けるために、Flink と Hologres の間で DECIMAL 型の精度が一貫していることを確認してください。詳細については、「よくある質問」をご参照ください。

  • ソーステーブルから完全および増分データを統合して消費するために initial モードを使用する場合、グローバルな順序は保証されません。下流システムが時間ベースの計算に依存している場合は、別の binlog のみの消費モードを使用してください。

binlog の有効化

新しいテーブルの場合

リアルタイムデータ読み取り機能はデフォルトで無効になっています。そのため、DDL ステートメントを使用して HoloWeb でテーブルを作成する際には、binlog.level および binlog.ttl パラメーターを設定する必要があります。以下に例を示します。

begin;
create table test_table(
  id int primary key, 
  title text not null, 
  body text);
call set_table_property('test_table', 'orientation', 'row');--test_table という名前の行指向テーブルを作成します。
call set_table_property('test_table', 'clustering_key', 'id');--id 列にクラスタリングキーを作成します。
call set_table_property('test_table', 'binlog.level', 'replica');--binlog 機能を有効にします。
call set_table_property('test_table', 'binlog.ttl', '86400');--binlog の生存時間 (TTL) を秒単位で設定します。
commit;

既存のテーブルの場合

HoloWeb で、次のステートメントを使用して既存のテーブルの Binlog を有効にし、Binlog TTL を設定できます。table_name は、Binlog を有効にしたいテーブルの名前です。

-- binlog 機能を有効にします。
begin;
call set_table_property('<table_name>', 'binlog.level', 'replica');
commit;

-- binlog の生存時間 (TTL) を秒単位で設定します。
begin;
call set_table_property('<table_name>', 'binlog.ttl', '2592000');
commit;

WITH パラメーター

VVR 11 では、一部の Hologres コネクタオプションの名前が変更または削除されましたが、VVR 8 との下位互換性は維持されています。ご利用のバージョンのパラメータードキュメントをご参照ください。

型マッピング

Flink と Hologres 間のデータ型マッピングについては、「Flink と Hologres 間のデータ型マッピング」をご参照ください。

説明

Hologres は、生成列を定義するために GENERATED ALWAYS AS 構文をサポートしています。例:

ds TIMESTAMP NOT NULL GENERATED ALWAYS AS (date_trunc('month', create_time)) STORED

生成列の NOT NULL 制約は、期待どおりにnull 許容フィールドにマッピングされます。生成列の値は Hologres によって計算されるため、Flink はこのフィールドに値を書き込みません。NOT NULL 制約が維持された場合、書き込み操作は Hologres クライアントでの検証に失敗します。この動作は、通常の列の NOT NULL 制約には影響しません (例:ds TIMESTAMP NOT NULL)。

例

ソーステーブルの例

Binlog ソーステーブル

CDC モード

このモードでは、ソースは binlog データを消費し、hg_binlog_event_type に基づいて各行に適切な Flink RowKind 型を自動的に設定するため、明示的な宣言は不要です。これらの型には、INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER が含まれます。これにより、MySQL や PostgreSQL などのデータベースの Change Data Capture (CDC) 機能と同様に、テーブルデータのミラーリングが可能になります。以下は、ソーステーブルの DDL の例です。

VVR 11+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='${secret_values.ak_id}',            --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。 
  'password'='${secret_values.ak_secret}',        
  'endpoint'='<yourEndpoint>',
  'source.binlog.change-log-mode'='ALL',  --INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER を含むすべての変更ログタイプを読み取ります。
  'retry-count'='10',                     --binlog 読み取りエラー時の再試行回数。
  'retry-sleep-step-ms'='5000',           --再試行間の増分バックオフ時間。最初の再試行は 5 秒、2 回目は 10 秒待機します。
  'source.binlog.batch-size'='512'        --binlog データを読み取るためのバッチサイズ。
);

VVR 8+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'cdcMode' = 'true',
  'sdkMode'='jdbc',
  'binlogMaxRetryTimes' = '10',     --binlog 読み取りエラー時の再試行回数。
  'binlogRetryIntervalMs' = '500',  --binlog 読み取りエラー後の再試行間隔 (ミリ秒)。
  'binlogBatchReadSize' = '100'     --binlog データを読み取るためのバッチサイズ。
);

非 CDC モード

このモードでは、ソースは消費した binlog データを通常の Flink データとして下流ノードに渡し、すべてのレコードを INSERT 型として扱います。ビジネス要件に基づいて、特定の hg_binlog_event_type を持つレコードを処理できます。以下は、ソーステーブルの DDL の例です。

VVR 11+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'source.binlog.change-log-mode'='ALL_AS_APPEND_ONLY',  --すべての変更ログタイプは INSERT 操作として扱われます。
  'retry-count'='10',                     --binlog 読み取りエラー時の再試行回数。
  'retry-sleep-step-ms'='5000',           --再試行間の増分バックオフ時間。最初の再試行は 5 秒、2 回目は 10 秒待機します。
  'source.binlog.batch-size'='512'        --binlog データを読み取るためのバッチサイズ。
);

VVR 8+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'binlogMaxRetryTimes' = '10',     --binlog 読み取りエラー時の再試行回数。
  'binlogRetryIntervalMs' = '500',  --binlog 読み取りエラー後の再試行間隔 (ミリ秒)。
  'binlogBatchReadSize' = '100'     --binlog データを読み取るためのバッチサイズ。
);

非 binlog ソーステーブル

VVR 11+

重要

VVR 11.1 以降、コネクタはデフォルトで binlog データを消費します。非 binlog ソーステーブルから読み取るには、'source.binlog' を 'false' に明示的に設定する必要があります。詳細については、「Binlog ソーステーブル」をご参照ください。

CREATE TEMPORARY TABLE hologres_source (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'source.binlog'='false'                      --binlog データを消費するかどうかを指定します。
);

VVR 8+

CREATE TEMPORARY TABLE hologres_source (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'sdkMode' = 'jdbc'
);

結果テーブル

CREATE TEMPORARY TABLE datagen_source(
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='datagen'
);
CREATE TEMPORARY TABLE hologres_sink (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       -- AK/SK キーの漏洩を防ぐために変数管理を使用します。
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>'
);
INSERT INTO hologres_sink SELECT * from datagen_source;

ディメンションテーブルの例

CREATE TEMPORARY TABLE datagen_source (
  a INT,
  b BIGINT,
  c STRING,
  proctime AS PROCTIME()
) WITH (
  'connector' = 'datagen'
);
CREATE TEMPORARY TABLE hologres_dim (
  a INT, 
  b VARCHAR, 
  c VARCHAR
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       -- キーの漏洩を防ぐために、AccessKey ペアには変数を使用します。
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>'
);
CREATE TEMPORARY TABLE blackhole_sink (
  a INT,
  b STRING
) WITH (
  'connector' = 'blackhole'
);
INSERT INTO blackhole_sink SELECT T.a,H.b
FROM datagen_source AS T JOIN hologres_dim FOR SYSTEM_TIME AS OF T.proctime AS H ON T.a = H.a;

高度な機能

全量と増分の統合取り込み

シナリオ

  • この機能は、プライマリキーを持つソーステーブルにのみ適用されます。CDC モードを使用する Hologres ソーステーブルに推奨されます。

  • Hologres では、必要に応じて Binlog を有効にできます。既にデータが含まれている既存のテーブルに対して Binlog を有効にすることができます。

コード例

VVR 11+

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'source.binlog.startup-mode' = 'INITIAL',   --すべての既存データを読み取り、その後 Binlog を増分的に消費します。
  'retry-count'='10',                         --Binlog データの読み取り中にエラーが発生した場合の再試行回数。
  'retry-sleep-step-ms'='5000',               --再試行間の増分待機時間。最初の再試行は 5 秒、2 回目は 10 秒待機します。
  'source.binlog.batch-size'='512'            --1 回のバッチで Binlog から読み取る行数。
  );
説明
  • source.binlog.startup-mode を INITIAL に設定すると、テーブルの初期完全読み取りを実行してから、Binlog の増分消費に切り替わります。

  • startTime パラメーターが直接設定されているか、起動 UI で開始時刻を選択して設定されている場合、それは優先され、binlogStartUpMode を timestamp モードに設定し、他のモード設定を上書きします。なぜなら、startTime パラメーターの方が優先度が高いためです。

VVR 8+

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'cdcMode' = 'true',
  'binlogStartUpMode' = 'initial', --すべての既存データを読み取り、その後 Binlog を増分的に消費します。
  'binlogMaxRetryTimes' = '10',     --Binlog データの読み取り中にエラーが発生した場合の再試行回数。
  'binlogRetryIntervalMs' = '500',  --Binlog 読み取りエラー後の再試行間隔 (ミリ秒)。
  'binlogBatchReadSize' = '100'     --1 回のバッチで Binlog から読み取る行数。
  );
説明
  • binlogStartUpMode を initial に設定すると、テーブルの初期完全読み取りを実行してから、Binlog の増分消費に切り替わります。

  • startTime パラメーターが直接設定されているか、起動 UI で開始時刻を選択して設定されている場合、それは優先され、binlogStartUpMode を timestamp モードに設定し、他のモード設定を上書きします。なぜなら、startTime パラメーターの方が優先度が高いためです。

プライマリキーの競合解決

コネクタは、書き込み時の重複プライマリキーを処理するための 3 つの戦略を提供します。

VVR 11+

sink.on-conflict-action パラメーターを設定することで、戦略を指定できます。

値

説明

INSERT_OR_IGNORE

最初に到着したレコードを保持し、後続の重複を無視します。

INSERT_OR_REPLACE

既存のレコードを新しいレコードで上書きします。

INSERT_OR_UPDATE (デフォルト)

結果テーブルで提供された列のみを更新し、既存のレコードの他の列は変更しません。

VVR 8+

mutatetype パラメーターを設定することで、戦略を指定できます。

値

説明

insertorignore (デフォルト)

最初に到着したレコードを保持し、後続の重複を無視します。

insertorreplace

既存のレコードを新しいレコードで上書きします。

insertorupdate

結果テーブルで提供された列のみを更新し、既存のレコードの他の列は変更しません。

例えば、テーブルに列 a、b、c、d があり、a がプライマリキーであるとします。結果テーブルが列 a と b のみを提供する場合、戦略を INSERT_OR_UPDATE に設定すると、列 b のみが更新されます。列 c と d は変更されません。
説明

ただし、結果テーブルから省略された物理テーブルの列は、null 許容でなければなりません。そうでない場合、書き込み操作は失敗します。

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

デフォルトでは、Hologres シンクは単一の非パーティション化テーブルにデータを書き込みます。パーティションテーブルにその親テーブルをターゲットとして書き込むには、以下のオプションを有効にする必要があります。

VVR 11+

コネクタが存在しない場合に子パーティションを自動的に作成できるようにするには、sink.create-missing-partition を true に設定します。

説明
  • VVR 11.1 以降のバージョンは、デフォルトでパーティションテーブルへの書き込みをサポートし、データを正しい子パーティションに自動的にルーティングします。

  • tablename パラメーターを親テーブルの名前に設定します。

  • 必要な子パーティションが存在せず、sink.create-missing-partition=true が設定されていない場合、書き込み操作は失敗します。

VVR 8+

  • データを対応する子パーティションに自動的にルーティングするには、partitionRouter を true に設定します。

  • コネクタが存在しない場合に子パーティションを自動的に作成できるようにするには、createparttable を true に設定します。

説明
  • tablename パラメーターを親テーブルの名前に設定します。

  • 必要な子パーティションが存在せず、createparttable=true が設定されていない場合、書き込み操作は失敗します。

ストリームのマージと部分更新

複数のストリームを単一の Hologres ワイドテーブルに書き込む際、コネクタは同じプライマリキーを持つレコードをマージします。部分更新は、行全体を置き換えるのではなく、変更された列のみを書き込むため、書き込みパフォーマンスとデータ整合性が向上します。

制限事項

  • ワイドテーブルにはプライマリキーが必要です。

  • 各データストリームには、プライマリキーを構成するすべての列が含まれている必要があります。

  • 列指向ストレージを使用するワイドテーブルの場合、高い秒間リクエスト数 (RPS) でストリームをマージすると、CPU 使用率が高くなる可能性があります。これを軽減するには、テーブルの列のディクショナリエンコーディングを無効にすることを検討してください。

例

2 つの Flink データストリームがあるとします。最初のストリームには列 a、b、c が含まれています。2 番目のストリームには列 a、d、e が含まれています。Hologres ワイドテーブル WIDE_TABLE には、列 a、b、c、d、e が含まれており、列 a がプライマリキーです。

VVR 11+

// source1 と source2 は既に定義されています。
CREATE TEMPORARY TABLE hologres_sink ( -- 列 a、b、c、d、e を宣言します。
  a BIGINT, 
  b STRING,
  c STRING,
  d STRING,
  e STRING,
  primary key(a) not enforced
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourWideTablename>',  -- Hologres ワイドテーブル。列 a、b、c、d、e を含みます。
  'username' = '${secret_values.ak_id}',
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'sink.on-conflict-action'='INSERT_OR_UPDATE',   -- プライマリキーに基づいて特定の列を更新します。
  'sink.delete-strategy'='IGNORE_DELETE',         -- リトラクションメッセージの処理戦略。IGNORE_DELETE は、削除操作が不要な追記専用またはアップサートストリームに適しています。
  'sink.partial-insert.enabled'='true'            -- 部分更新を有効にします。INSERT ステートメントで指定された列のみがコネクタに送信されます。
);

BEGIN STATEMENT SET;
INSERT INTO hologres_sink(a,b,c) select * from source1;  -- 列 a、b、c のみが挿入されることを宣言します。
INSERT INTO hologres_sink(a,d,e) select * from source2;  -- 列 a、d、e のみが挿入されることを宣言します。
END;

VVR 8+

// source1 と source2 は既に定義されています。
CREATE TEMPORARY TABLE hologres_sink ( -- 列 a、b、c、d、e を宣言します。
  a BIGINT, 
  b STRING,
  c STRING,
  d STRING,
  e STRING,
  primary key(a) not enforced
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourWideTablename>',  -- Hologres ワイドテーブル。列 a、b、c、d、e を含みます。
  'username' = '${secret_values.ak_id}',
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'mutatetype'='insertorupdate',    -- プライマリキーに基づいて特定の列を更新します。
  'ignoredelete'='true',            -- リトラクションメッセージによって生成された DELETE リクエストを無視します。
  'partial-insert.enabled'='true'   -- 部分更新を有効にして、INSERT ステートメントで宣言された列のみを更新します。
);

BEGIN STATEMENT SET;
INSERT INTO hologres_sink(a,b,c) select * from source1;  -- 列 a、b、c のみが挿入されることを宣言します。
INSERT INTO hologres_sink(a,d,e) select * from source2;  -- 列 a、d、e のみが挿入されることを宣言します。
END;
説明

ignoredelete を true に設定すると、リトラクションメッセージによって生成される Delete リクエストが無視されます。VVR 8.0.8 以降では、sink.delete-strategy を使用して、部分更新のさまざまな戦略を設定することを推奨します。

パーティションテーブルからの binlog 消費 (ベータ版)

Hologres コネクタは、物理パーティションテーブルと論理パーティションテーブルの両方からの binlog 消費をサポートしています。両者の違いについては、「CREATE LOGICAL PARTITION TABLE」をご参照ください。

物理パーティションテーブルからの binlog 消費

Hologres コネクタは、パーティションテーブルから binlog を消費し、単一のジョブ内で新しいパーティションを動的に監視できます。これにより、リアルタイムのデータ処理効率と使いやすさが大幅に向上します。

注意事項

  • この機能は、JDBC モードの binlog ソーステーブルでのみ利用可能です。VVR 8.0.11 以降と Hologres インスタンスのバージョン 2.1.27 以降が必要です。

  • パーティション名は、動的パーティション分割の命名規則に従う必要があります:{parent_table}_{partition_value}。準拠していないパーティションは消費されない可能性があります。

    重要
    • DYNAMIC モードの場合、VVR バージョン 8.0.11 は、- デリミタ (YYYY-MM-DD など) を持つパーティション列をサポートしていません。

    • VVR 11.1 以降では、カスタム形式を使用するパーティション列からデータを消費できます。

    • この制限は、パーティションテーブルへの書き込みには適用されません。

  • Flink で Hologres ソーステーブルを宣言する際には、そのパーティション列を含める必要があります。

  • DYNAMIC モードでは、パーティションテーブルで動的パーティション分割が有効になっている必要があります。さらに、パーティション事前作成パラメーター auto_partitioning.num_precreate は 1 より大きい必要があります。そうでない場合、ジョブは最新のパーティションを消費しようとすると例外をスローします。

  • DYNAMIC モードでは、新しいパーティションが追加された後、コネクタは古いパーティションからの後続のデータ変更を消費しなくなります。

例

モード

特徴

説明

DYNAMIC

動的パーティション消費

新しいパーティションを時系列順に自動的に監視および消費します。このモードは、リアルタイムのデータストリーミングシナリオに適しています。

STATIC

静的パーティション消費

既存のパーティション (または手動で指定されたもの) のみを消費し、新しいパーティションを自動的に検出しません。このモードは、固定範囲内の既存データの処理に適しています。

Dynamic モード

VVR 11+

Hologres パーティションテーブルが以下の DDL で作成され、binlog と動的パーティション分割が有効になっているとします。

CREATE TABLE "test_message_src1" (
    id int,
    title text,
    body text,
    dt text,
    PRIMARY KEY (id, dt)
)
PARTITION BY LIST (dt) WITH (
    binlog_level = 'replica', 
    auto_partitioning_enable =  'true',   -- 動的パーティション分割を有効にします。
    auto_partitioning_time_unit = 'DAY',  -- パーティションは毎日作成されます。パーティション名の例:test_message_src1_20250512, test_message_src1_20250513。
    auto_partitioning_num_precreate = '2' -- 2 つのパーティションを事前に作成します。
);
-- 既存のパーティションテーブルの場合、ALTER TABLE を使用して動的パーティション分割を有効にすることもできます。

Flink で、次の SQL ステートメントを使用して、パーティションテーブル test_message_src1 を DYNAMIC モードで消費します。

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  dt VARCHAR  -- Hologres パーティションテーブルのパーティション列。
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src1',  -- 動的パーティション分割が有効になっている親テーブル。
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'source.binlog.partition-binlog-mode' = 'DYNAMIC', -- 最新のパーティションを動的に監視します。
  'source.binlog.startup-mode' = 'initial'           -- すべての既存データを消費し、その後 binlog から増分データを消費します。
);

VVR 8.0.11

Hologres パーティションテーブルが以下の DDL で作成され、binlog と動的パーティション分割が有効になっているとします。

CREATE TABLE "test_message_src1" (
    id int,
    title text,
    body text,
    dt text,
    PRIMARY KEY (id, dt)
)
PARTITION BY LIST (dt) WITH (
    binlog_level = 'replica', 
    auto_partitioning_enable =  'true',   -- 動的パーティション分割を有効にします。
    auto_partitioning_time_unit = 'DAY',  -- パーティションは毎日作成されます。パーティション名の例:test_message_src1_20241027, test_message_src1_20241028。
    auto_partitioning_num_precreate = '2' -- 2 つのパーティションを事前に作成します。
);

-- 既存のパーティションテーブルの場合、ALTER TABLE を使用して動的パーティション分割を有効にすることもできます。

Flink で、次の SQL ステートメントを使用して、パーティションテーブル test_message_src1 から DYNAMIC モードでデータを消費します。

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  dt VARCHAR  -- Hologres パーティションテーブルのパーティション列。
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src1',  -- 動的パーティション分割が有効になっている親テーブル。
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'binlog' = 'true',
  'partition-binlog.mode' = 'DYNAMIC',  -- 最新のパーティションを動的に監視します。
  'binlogstartUpMode' = 'initial',      -- すべての既存データを消費し、その後 binlog から増分データを消費します。
  'sdkMode' = 'jdbc_fixed'              -- 接続制限の問題を回避するためにこのモードを使用します。
);

Static モード

VVR 11+

Hologres パーティションテーブルが以下の DDL で作成され、binlog が有効になっているとします。

CREATE TABLE test_message_src2 (
    id int,
    title text,
    body text,
    color text,
    PRIMARY KEY (id, color)
)
PARTITION BY LIST (color) WITH (
    binlog_level = 'replica'
);
create table test_message_src2_red partition of test_message_src2 for values in ('red');
create table test_message_src2_blue partition of test_message_src2 for values in ('blue');
create table test_message_src2_green partition of test_message_src2 for values in ('green');
create table test_message_src2_black partition of test_message_src2 for values in ('black');

Flink で、次の SQL ステートメントを使用して、パーティションテーブル test_message_src2 を STATIC モードで消費します。

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  color VARCHAR  -- Hologres パーティションテーブルのパーティション列。
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src2',  -- パーティションテーブル。
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'source.binlog.partition-binlog-mode' = 'STATIC', -- 固定されたパーティションセットを消費します。
  'source.binlog.partition-values-to-read' = 'red,blue,green',  -- 指定された 3 つのパーティションのみを消費します。「black」パーティションは消費されません。新しいパーティションも消費されません。このオプションが設定されていない場合、ジョブは親テーブルのすべてのパーティションを消費します。
  'source.binlog.startup-mode' = 'initial'  -- すべての既存データを消費し、その後 binlog から増分データを消費します。
);

VVR 8.0.11

Hologres パーティションテーブルが以下の DDL で作成され、binlog が有効になっているとします。

CREATE TABLE test_message_src2 (
    id int,
    title text,
    body text,
    color text,
    PRIMARY KEY (id, color)
)
PARTITION BY LIST (color) WITH (
    binlog_level = 'replica'
);
create table test_message_src2_red partition of test_message_src2 for values in ('red');
create table test_message_src2_blue partition of test_message_src2 for values in ('blue');
create table test_message_src2_green partition of test_message_src2 for values in ('green');
create table test_message_src2_black partition of test_message_src2 for values in ('black');

Flink で、次の SQL ステートメントを使用して、パーティションテーブル test_message_src2 から STATIC モードでデータを消費します。

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  color VARCHAR  -- Hologres パーティションテーブルのパーティション列。
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src2',  -- パーティションテーブル。
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'binlog' = 'true',
  'partition-binlog.mode' = 'STATIC', -- 固定されたパーティションセットを消費します。
  'partition-values-to-read' = 'red,blue,green',  -- 指定された 3 つのパーティションのみを消費します。「black」パーティションは消費されません。新しいパーティションも消費されません。このオプションが設定されていない場合、ジョブは親テーブルのすべてのパーティションを消費します。
  'binlogstartUpMode' = 'initial',  -- すべての既存データを消費し、その後 binlog から増分データを消費します。
  'sdkMode' = 'jdbc_fixed' -- 接続制限の問題を回避するためにこのモードを使用します。
);

論理パーティションテーブルからの binlog 消費

Hologres コネクタは、論理パーティションテーブルからの binlog 消費をサポートしており、オプションで消費するパーティションを指定できます。

注意事項

  • 論理パーティションテーブルの特定のパーティションから binlog を消費するには、VVR 11.0.0 以降と Hologres インスタンスのバージョン V3.1 以降が必要です。

  • 論理パーティションテーブルのすべてのパーティションから binlog を消費する方法は、非パーティション化テーブルと同じアプローチに従います (ソーステーブル)。

例

パラメーター

説明

例

source.binlog.logical-partition-filter-column-names

消費するパーティションを指定するパーティション列名。列名は二重引用符 (") で囲みます。複数の列名はカンマ (,) で区切ります。列名に二重引用符が含まれる場合は、追加の二重引用符でエスケープします。

'source.binlog.logical-partition-filter-column-names'='"Pt","id"'

Pt と id の 2 つのパーティション列が使用されます。

source.binlog.logical-partition-filter-column-values

消費するパーティションを指定するパーティション値。パーティションは、各パーティション列に 1 つずつ、値のセットで指定されます。各値は二重引用符 (") で囲みます。同じパーティションの値はカンマ (,) で区切ります。パーティションはセミコロン (;) で区切ります。値に二重引用符が含まれる場合は、追加の二重引用符でエスケープします。

'source.binlog.logical-partition-filter-column-values'='"20240910","0";"special""value","9"'

これにより、消費する 2 つのパーティションが指定されます。最初のパーティション値は (20240910, 0) で、2 番目は (special"value, 9) です。

Hologres で以下のテーブルを作成したとします。

CREATE TABLE holo_table (
    id int not null,
    name text,
    age numeric(18,4),
    "Pt" text,
    primary key(id, "Pt")
)
LOGICAL PARTITION BY LIST ("Pt", id)
WITH (
    binlog_level ='replica'
);

このテーブルの binlog を Flink で消費するには:

CREATE TEMPORARY TABLE test_src_binlog_table(
  id INTEGER,
  name VARCHAR,
  age decimal(18,4),
  `Pt` VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='holo_table',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'source.binlog'='true',
  'source.binlog.logical-partition-filter-column-names'='"Pt","id"',
  'source.binlog.logical-partition-filter-column-values'='<yourPartitionColumnValues>',
  'source.binlog.change-log-mode'='ALL',  --INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER を含むすべての変更ログタイプを読み取ります。
  'retry-count'='10',                     -- binlog 読み取りエラーの再試行回数。
  'retry-sleep-step-ms'='5000',           --再試行間の増分バックオフ時間。最初の再試行は 5 秒、2 回目は 10 秒待機します。
  'source.binlog.batch-size'='512'        --1 回のバッチで binlog から読み取る行数。
);

DataStream API

重要

DataStream API を使用して Hologres から読み書きするには、対応する DataStream コネクタを設定します (DataStream コネクタの使用方法)。Hologres DataStream コネクタは Maven Central で入手できます。ローカルでのデバッグには、Uber JAR を使用します (コネクタを含むジョブをローカルで実行およびデバッグする)。

Hologres ソーステーブル

Binlog ソーステーブル

VVR は、Hologres binlog データを読み取るための HologresBinlogSource クラスを提供します。次の例は、Hologres binlog ソースを構築する方法を示しています。

VVR 11.3+

重要

VVR 11.1.2 以降、JDBCOptions および startTimeMs パラメーターは HologresBinlogSource コンストラクターから削除されました。VVR 11.3 以降、List<Subscribe.BinlogFilter> パラメーターが追加されました。VVR 11 以降を使用する場合は、VVR 11.3 以降を使用することを推奨します。

public class Sample {                                                                                                                                                                          
        public static void main(String[] args) throws Exception {
            final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            // 読み取るテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
            TableSchema schema = TableSchema.builder()
                    .field("a", DataTypes.INT())
                    .field("b", DataTypes.STRING())
                    .field("c", DataTypes.TIMESTAMP())
                    .build();

            // 読み取るテーブルの名前。
            String sourceTableName = "sourceTableName";

            // Hologres 接続のパラメーター。
            Configuration config = new Configuration();
            config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
            config.setString(HologresConfigs.USERNAME, "yourUserName");
            config.setString(HologresConfigs.PASSWORD, "yourPassword");
            config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
            config.setString(HologresConfigs.TABLE, sourceTableName);
            config.set(HologresConfigs.BINLOG, true);
            config.set(HologresConfigs.BINLOG_CHANGE_LOG_MODE, BinlogChangeLogMode.ALL);
            // Hologres binlog ソースを構築します。
            HologresBinlogSource source = new HologresBinlogSource(
                    new HologresConnectionParam(config),
                    schema,
                    config,
                    StartupMode.INITIAL,
                    sourceTableName,
                    "",
                    Collections.emptyList(),
                    -1,
                    Collections.emptySet(),
                    Collections.emptyList()
            );
            env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
            env.execute();
        }
  }

VVR 8.0.11+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 読み取るテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Hologres 接続のパラメーター。
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // JDBCOptions を構築します。
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // Hologres binlog ソースを構築します。
        long startTimeMs = 0;
        HologresBinlogSource source = new HologresBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.INITIAL,
                "",
                "",
                -1,
                Collections.emptySet(),
                new ArrayList<>()
        );
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}

VVR 8.0.7+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 読み取るテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Hologres 接続のパラメーター。
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // JDBCOptions を構築します。
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // Hologres binlog ソースを構築します。
        long startTimeMs = 0;
        HologresBinlogSource source = new HologresBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.INITIAL,
                "",
                "",
                -1,
                Collections.emptySet()
        );
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}

VVR 6.0.7+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 読み取るテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .build();
         // Hologres 接続のパラメーター。
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // JDBCOptions を構築します。
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // デフォルトのスロット名を設定または作成します。
        config.setString(HologresBinlogConfigs.JDBC_BINLOG_SLOT_NAME, HoloBinlogUtil.getOrCreateDefaultSlotForJDBCBinlog(jdbcOptions));

        boolean cdcMode = config.get(HologresBinlogConfigs.BINLOG_CDC_MODE) && config.get(HologresBinlogConfigs.OPTIONAL_BINLOG);
        // JDBCBinlogRecordConverter を構築します。
        JDBCBinlogRecordConverter recordConverter = new JDBCBinlogRecordConverter(
                jdbcOptions.getTable(),
                schema,
                new HologresConnectionParam(config),
                cdcMode,
                Collections.emptySet());
        
        // Hologres binlog ソースを構築します。
        long startTimeMs = 0;
        HologresJDBCBinlogSource source = new HologresJDBCBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.TIMESTAMP,
                recordConverter,
                "",
                -1);
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}
重要

Flink エンジンバージョンが 8.0.5 より前、または Hologres が 2.1 より前の場合、ユーザーがスーパーユーザーであるか、レプリケーションロールを持っていることを確認してください (Hologres の権限の問題)。

非 binlog ソーステーブル

VVR は、Hologres テーブルからデータを読み取るための RichInputFormat の実装である HologresBulkreadInputFormat クラスを提供します。次の例は、Hologres ソースを構築する方法を示しています。

public class Sample {
    public static void main(String[] args) throws Exception {
        // Java DataStream API を設定します
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 読み取るテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Hologres 接続のパラメーター。
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        // JDBCOptions を構築します。
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        HologresBulkreadInputFormat inputFormat = new HologresBulkreadInputFormat(
                new HologresConnectionParam(config),
                jdbcOptions,
                schema,
                "",
                -1);
        TypeInformation<RowData> typeInfo = InternalTypeInfo.of(schema.toRowDataType().getLogicalType());
        env.addSource(new InputFormatSourceFunction<>(inputFormat, typeInfo)).returns(typeInfo).print();
        env.execute();
    }
}

Maven 依存関係

Hologres DataStream コネクタは Maven Central で入手できます。

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-hologres</artifactId>
    <version>${vvr-version}</version>
</dependency>

Hologres 結果テーブル

VVR 11+

public class Sample {
      public static void main(String[] args) throws Exception {
          final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
          // 書き込むテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
          TableSchema tableSchema = TableSchema.builder()
                  .field("a", DataTypes.INT().notNull())
                  .field("b", DataTypes.STRING())
                  .primaryKey("a")
                  .build();
          // Hologres 接続のパラメーター。
          Configuration config = new Configuration();
          config.set(HologresConfigs.ENDPOINT, "yourEndpoint");
          config.set(HologresConfigs.USERNAME, "yourUserName");
          config.set(HologresConfigs.PASSWORD, "yourPassword");
          config.set(HologresConfigs.DATABASE, "yourDatabaseName");
          config.set(HologresConfigs.TABLE, "yourTableName");
          HologresConnectionParam connectionParam = new HologresConnectionParam(config);
          HologresTableSchema hologresTableSchema =
                  HologresTableSchema.get(connectionParam.getJDBCOptions());
          // 結果テーブルに書き込む列のインデックス。
          Integer[] targetColumnIndexes = {0, 1};
          // Hologres シンクを構築します。
          HologresSinkFunction sinkFunction =
                  new HologresSinkFunction(
                          connectionParam, tableSchema, targetColumnIndexes, hologresTableSchema);
          TypeInformation<RowData> typeInfo = InternalTypeInfo.of(tableSchema.toRowDataType().getLogicalType());
          env.fromElements((RowData) GenericRowData.of(101, StringData.fromString("name"))).returns(typeInfo).addSink(sinkFunction);
          env.execute();
      }
  }

VVR 8+

public class Sample {
    public static void main(String[] args) throws Exception {
        // Java DataStream API を設定します
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 書き込むテーブルのスキーマを初期化します。スキーマは Hologres テーブルスキーマのフィールドと一致する必要があります。フィールドのサブセットを定義できます。
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .build();
        // Hologres 接続のパラメーター。
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        HologresConnectionParam hologresConnectionParam = new HologresConnectionParam(config);
        
         // RowData としてデータを書き込むための Hologres ライターを構築します。
        AbstractHologresWriter<RowData> hologresWriter = HologresJDBCWriter.createRowDataWriter(
                hologresConnectionParam, 
                schema, 
                HologresTableSchema.get(hologresConnectionParam), 
                new Integer[0]);
        // Hologres シンクを構築します。
        HologresSinkFunction sinkFunction = new HologresSinkFunction(hologresConnectionParam, hologresWriter);
        TypeInformation<RowData> typeInfo = InternalTypeInfo.of(schema.toRowDataType().getLogicalType());
        env.fromElements((RowData) GenericRowData.of(101, StringData.fromString("name"))).returns(typeInfo).addSink(sinkFunction);
        env.execute();
    }
}

メタデータ列

VVR 8.0.11 以降は、binlog ソーステーブルでメタデータ列をサポートしています。hg_binlog_event_type などの binlog フィールドをメタデータ列として宣言すると、ソースデータベース名、テーブル名、変更タイプ、イベントタイムスタンプにアクセスして、カスタムロジック (例:DELETE イベントのフィルタリング) を実装できます。

パラメーター

型

説明

db_name

STRING NOT NULL

行を含むデータベースの名前。

table_name

STRING NOT NULL

行を含むテーブルの名前。

hg_binlog_lsn

BIGINT NOT NULL

binlog シーケンス番号のシステムフィールド。値はシャード内で単調に増加しますが、連続していません。シャード間での一意性と順序は保証されません。

hg_binlog_timestamp_us

BIGINT NOT NULL

データベースでの変更イベントのタイムスタンプ (マイクロ秒単位)。

hg_binlog_event_type

BIGINT NOT NULL

行の変更タイプ。有効な値は次のとおりです:

  • 5: INSERT メッセージ。

  • 2: DELETE メッセージ。

  • 3: UPDATE 操作前の行イメージ。

  • 7: UPDATE 操作後の行イメージ。

hg_shard_id

INT NOT NULL

行を含むデータシャードの ID (テーブルグループとシャード)。

DDL ステートメントでは、<meta_column_name> <datatype> METADATA VIRTUAL を使用してメタデータ列を宣言できます。以下に例を示します:

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn bigint METADATA VIRTUAL
  hg_binlog_event_type bigint METADATA VIRTUAL
  hg_binlog_timestamp_us bigint METADATA VIRTUAL
  hg_shard_id int METADATA VIRTUAL
  db_name string METADATA VIRTUAL
  table_name string METADATA VIRTUAL
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  ...
  );

よくある質問

参考