Hologres は、専用モードの Blink で実行される Realtime Compute for Apache Flink とシームレスに統合されます。コネクタを使用することで、データストリームを Hologres シンクテーブルに書き込み、すぐにデータをクエリできます。本トピックでは、専用モードの Blink で実行されるジョブから Hologres シンクテーブルにデータを書き込む方法について説明します。
制限事項
専用モードの Blink は、バージョンによって開発構文が異なります。開始する前に、お使いの専用モードの Blink のバージョンを特定し、対応する例を使用してください。
接続エラーを回避するために、Realtime Compute for Apache Flink のデプロイメントと Hologres のインスタンスが同じリージョンにあることを確認してください。
バージョン 3.6 より前の専用モードの Blink には、Hologres コネクタが組み込まれていません。Hologres にリアルタイムでデータを書き込むには、JAR ファイルを参照する必要があります。サポートが必要な場合は、「アップグレード準備中の一般的なエラー」をご参照いただくか、Hologres DingTalk グループに参加してお問い合わせください。詳細については、「オンラインサポートの利用方法」をご参照ください。
説明ジョブを実行するには、専用モードの Blink をバージョン 3.6 以降にアップグレードすることを推奨します。
Blink 専用モード 3.7 は、Hologres パーティションテーブルの自動作成をサポートしていますが、ジョブで
createparttable='true'を設定する必要があります。パーティションテーブルを使用する場合の考慮事項は次のとおりです。Hologres は現在、リストパーティショニングのみをサポートしています。
パーティションテーブルを作成する際には、パーティションキー列を明示的に指定する必要があります。現在、パーティションキー列は text と int4 のデータ型のみをサポートしており、その値には、たとえば
2020-09-12のようにハイフン (-) を含めることはできません。パーティションテーブルにプライマリキーが定義されている場合、パーティションキー列はプライマリキーの一部である必要があります。
子パーティションテーブルを作成する際、そのパーティションキー列の値は固定値である必要があります。
子パーティションテーブルに書き込まれるデータのパーティションキー列の値は、子テーブルの作成時に定義された値と完全に一致する必要があります。一致しない場合、エラーが発生します。
DEFAULT パーティション機能は現在サポートされていません。
宛先の Hologres テーブルにプライマリキーがある場合、デフォルトのリアルタイム書き込みセマンティクスでは、そのキーに基づいてレコードは更新されません。後で重複するプライマリキーを持つデータをインポートした場合、新しいデータは破棄されます。
Hologres はデータを非同期に書き込みます。ジョブで例外が発生した場合にフェイルオーバーがトリガーされるように、ジョブに
blink.checkpoint.fail_on_checkpoint_error=true設定を追加する必要があります。このパラメーターは、Blink 3.7.6 以降では不要です。
DDL 構文
次の DDL ステートメントは、Hologres シンクテーブルを作成します。
create table Hologres_sink(
name varchar,
age BIGINT,
birthday BIGINT
) with (
type='hologres',
dbname='<yourDbname>', -- Hologres データベースの名前。
tablename='<yourTablename>', -- データを受け取る Hologres テーブルの名前。
username='<yourUsername>', -- Alibaba Cloud アカウントの AccessKey ID。
password='<yourPassword>', -- Alibaba Cloud アカウントの AccessKey Secret。
endpoint='<yourEndpoint>'); -- Hologres インスタンスの VPC エンドポイント。WITH パラメーター
パラメーター | 説明 | 例 |
type | シンクテーブルのタイプ。 | hologres |
endpoint | Hologres インスタンスの VPC エンドポイント。 Hologres コンソールにログインし、インスタンス詳細ページの ネットワーク情報 セクションでエンドポイントを確認できます。エンドポイントにはポート番号を含め、「IP:ポート」の形式に従う必要があります。 | demo-cn-hangzhou-vpc.hologres.aliyuncs.com:80 |
username | AccessKey ID [AccessKey Management] ページで AccessKey ID を取得できます。 | xxxxm3FMWaxxxx |
password | AccessKey Secret [AccessKey Management] ページで AccessKey Secret を取得できます。 | xxxxm355fffaxxxx |
dbname | Hologres データベースの名前です。 | Holodb |
tablename | Hologres データベース内のテーブルの名前です。 | blink_test |
arraydelimiter | Hologres シンクは、この区切り文字を使用して STRING フィールドを配列に分割してから、その配列を Hologres にインポートします。 デフォルト値は\u0002です。 | \u0002 |
mutatetype | データ書き込みモード。詳細については、「Hologres シンクテーブル」をご参照ください。 デフォルト値は insertorignore です。 | insertorignore |
ignoredelete | リトラクションメッセージを無視するかどうかを指定します。
説明 このパラメーターはストリーミングジョブに対してのみ有効です。 デフォルト値は false です。 Flink の | フォールス |
partitionrouter | パーティションテーブルにデータを書き込むかどうかを指定します。
デフォルト値は false です。 | 偽 |
createparttable | パーティションテーブルに書き込む場合に、パーティション値に基づいて子パーティションテーブルを自動的に作成するかどうかを指定します。この機能は、専用モードの Blink バージョン 3.7 以降でサポートされています。 デフォルト値は false です。 重要 この機能は注意して使用してください。パーティション値にダーティデータが含まれていないことを確認してください。ダーティデータが含まれていると、不正なパーティションテーブルが作成される可能性があります。 | 無効 |
arraydelimiter、mutatetype、ignoredelete、partitionrouter、および createparttable パラメーターは、DDL 文の例には含まれていません。アプリケーションでこれらのパラメーターを使用する必要がある場合は、この表の説明に従って追加してください。
標準 Hologres シンクテーブルへの書き込み
Hologres にテーブルを作成します。
データを受け取るためのテーブルを Hologres に作成します。以下は SQL ステートメントの例です。
create table blink_test (a int, b text, c text, d float8, e bigint);Realtime Compute for Apache Flink ジョブを作成します。
ジョブを作成します。
専用モードの Blink バージョン 3.6 以降には、Hologres データソースのサポートが組み込まれています。このデータソースを直接使用できます。以下は SQL ステートメントの例です。
create table randomSource (a int, b VARCHAR, c VARCHAR, d DOUBLE, e BIGINT) with (type = 'random'); create table test ( a int, b VARCHAR, c VARCHAR, PRIMARY KEY (a) ) with ( type = 'hologres', `endpoint` = '$ip:$port', -- Hologres インスタンスの VPC エンドポイントとポート番号。 `username` = 'Alibaba Cloud アカウントの AccessKey ID。', `password` = 'Alibaba Cloud アカウントの AccessKey Secret。', `dbname` = 'Hologres データベースの名前。', `tablename` = 'blink_test'-- データを受け取る Hologres テーブルの名前。 ); insert into test select a,b,c from randomSource;
ジョブを公開します。
ジョブを作成した後、[構文チェック] をクリックします。成功しました ステータスは、構文が正しいことを示します。
保存をクリックしてジョブを保存します。
[公開] をクリックして、ジョブを本番環境にデプロイします。 ビジネス要件に合わせてデプロイ設定を構成します。 [新しいバージョンを公開] をクリックして、デプロイプロセスを開始します。 [初期リソース] ステップで、リソース割り当て方法として [最終実行に基づく自動チューニング]、[システム割り当て]、または [手動リソース設定] のいずれかを選択します。 選択後、[次へ] をクリックします。 または、[データチェックをスキップ] をクリックして、リソース設定ステップに直接進むこともできます。
ジョブを開始します。
ジョブを本番環境に公開した後、手動で開始する必要があります。
[開発プラットフォーム] ページの上部ナビゲーションバーの右側にある [管理] をクリックします。管理ページで、目的のジョブを選択し、右上隅の 起動 をクリックします。
Hologres 内のデータをリアルタイムでクエリする
Hologres の宛先テーブルをクエリして、書き込まれたデータをリアルタイムで表示します。以下はクエリの例です。
select * from blink_test;
ワイドテーブルのマージと更新
このセクションでは、複数のストリームからのデータを単一の Hologres ワイドテーブルに書き込むという、一般的なユースケースについて説明します。
Hologres に、列 A を主キーとする A、B、C、D、E の列を持つ WIDE_TABLE という名前のワイドテーブルがあり、Flink では、1 つのストリームには列 A、B、C のデータが含まれ、もう 1 つのストリームには列 A、D、E のデータが含まれているとします。
Flink SQL を使用して、2 つの Hologres シンクテーブルを宣言します。1 つのテーブルには列 A、B、C を、もう 1 つのテーブルには列 A、D、E を指定します。両方のテーブルを Hologres の WIDE_TABLE テーブルにマッピングします。
両方のシンクテーブルのmutatetypeパラメーターをinsertorupdateに設定します。
両方のシンクテーブルで、ignoredelete パラメーターを true に設定します。これにより、リトラクションメッセージが
DELETEリクエストを生成するのを防ぎます。各ストリームのデータを、対応するシンクテーブルに挿入します。
このシナリオには、次の制限事項があります:
ワイドテーブルにはプライマリキーが必要です。
各ストリームには、すべてのプライマリキー列を含める必要があります。
高い RPS でカラム指向のワイドテーブルにデータをマージすると、CPU 使用率が高くなる可能性があります。 テーブル内の列の ディクショナリエンコーディング を無効にすることをお勧めします。
パーティション化された Hologres シンクテーブルへの書き込み
Hologres では、HoloHub API を呼び出して、親パーティションテーブルに直接データを書き込むことができます。データは、その後、正しい子パーティションテーブルに自動的にルーティングされます。詳細については、「HoloHub API」をご参照ください。
適用される制限事項は次のとおりです:
Hologres は現在、リストパーティショニングのみをサポートしています。
パーティション化テーブルを作成する際には、パーティションキー列を明示的に指定する必要があります。パーティションキー列のデータ型は text または int4 のみです。
プライマリキーが定義されている場合、パーティションキー列はプライマリキーの一部である必要があります。
子パーティションテーブルを作成する際、そのパーティションキー列の値は固定値である必要があります。
子パーティションテーブルに書き込まれるデータのパーティションキー列の値は、子テーブルの作成時に定義された値と完全に一致する必要があります。一致しない場合、エラーが発生します。
Hologres は現在、デフォルトパーティションをサポートしていません。
Hologres にパーティションテーブルを作成します。
データを受け取るためのパーティションテーブルを Hologres に作成し、対応する子パーティションテーブルを作成します。以下は SQL ステートメントの例です。
-- 親テーブル test_message とその子パーティションテーブルを作成します。 drop table if exists test_message; begin; create table test_message ( "bizdate" text NOT NULL, "tag" text NOT NULL, "id" int4 NOT NULL, "title" text NOT NULL, "body" text, PRIMARY KEY (bizdate,tag,id) ) PARTITION BY LIST (bizdate); commit;説明コマンドを実行する際、
${bizdate}パラメーターを実際の値に置き換えてください。パーティションの自動作成をサポートしているのは、バージョン 3.7 以降の専用モードの Blink のみです。以前のバージョンを使用している場合は、事前に Hologres で子パーティションテーブルを作成する必要があります。そうしないと、データのインポートに失敗します。
Blink で排他モードでジョブを作成する
以下は、専用モードの Blink でジョブを作成するためのステートメントの例です。
説明以下の例は、Blink exclusive mode 3.7 以降に適用されます。Blink exclusive mode 3.7 より前のバージョンを使用している場合は、バージョン 3.7 以降にアップグレードするか、
createparttable = 'true'の設定を削除してください。create table test_message_src( tag VARCHAR, id INTEGER, title VARCHAR, body VARCHAR ) with ( type = 'random', `interval` = '10', `count` = '100' ); create table test_message_sink ( bizdate VARCHAR, tag VARCHAR, id INTEGER, title VARCHAR, body VARCHAR ) with ( type = 'hologres', `endpoint` = '$ip:$port', -- Hologres インスタンスの VPC エンドポイント。 `username` ='<AccessID>', -- Alibaba Cloud アカウントの AccessKey ID。 `password` = '<AccessKey>', -- Alibaba Cloud アカウントの AccessKey Secret。 `dbname` = '<DBname>', -- Hologres データベースの名前。 `tablename` = '<Tablename>', -- Hologres データベース内のテーブルの名前。 `partitionrouter` = 'true', -- Hologres のパーティションテーブルにデータを書き込みます。 `createparttable` = 'true' -- Hologres で子パーティションテーブルを自動的に作成します。 ); insert into test_message_sink select "20200327",* from test_message_src; insert into test_message_sink select "20200328",* from test_message_src;ジョブを公開して開始します。
詳細については、「標準の Hologres シンクテーブルにデータを書き込む」セクションの「ジョブの公開」および「ジョブの開始」の手順をご参照ください。
Hologres のデータをリアルタイムでクエリする。
Hologres の宛先テーブルをクエリして、書き込まれたデータをリアルタイムで表示します。以下はクエリの例です。
select * from test_message; select * from test_message where bizdate = '20200327';
データ型のマッピング
専用モードの Blink と Hologres 間のデータ型のマッピングについては、「データ型の概要」をご参照ください。