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 コンソールにログインし、インスタンス詳細ページの Network Information セクションでエンドポイントを確認します。エンドポイントはポート番号を含め、ip:port 形式である必要があります。 |
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 の |
false |
|
partitionrouter |
パーティションテーブルにデータを書き込むかどうかを指定します。
デフォルト値は false です。 |
false |
|
createparttable |
パーティションテーブルに書き込む場合に、パーティション値に基づいて子パーティションテーブルを自動的に作成するかどうかを指定します。この機能は、専用モードの Blink バージョン 3.7 以降でサポートされています。 デフォルト値は false です。 重要
この機能は注意して使用してください。パーティション値にダーティデータが含まれていないことを確認してください。ダーティデータがあると、不正なパーティションテーブルが作成される可能性があります。 |
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;
-
-
ジョブを公開します。
-
ジョブを作成した後、[Syntax Check] をクリックします。成功しました ステータスは、構文が正しいことを示します。
-
保存 をクリックしてジョブを保存します。
-
[Publish] をクリックして、ジョブを本番環境にデプロイします。ビジネス要件に基づいてデプロイ設定を構成します。[Publish New Version] をクリックしてデプロイプロセスを開始します。[Initial Resources] ステップで、リソース割り当て方法 ([Auto-tuning based on last run]、[System allocation]、または [Manual resource configuration]) を選択します。選択後、[Next] をクリックします。[Skip Data Check] をクリックして、リソース設定ステップに直接進むこともできます。
-
-
ジョブを開始します。
ジョブを本番環境に公開した後、手動で開始する必要があります。
[Development Platform] ページの上部メニューで、右側の [Administration] をクリックします。[Administration] ページで、目的のジョブを選択し、右上隅の 起動 をクリックします。
-
Hologres のデータをリアルタイムでクエリします。
Hologres の宛先テーブルをクエリして、書き込まれたデータをリアルタイムで表示します。以下にクエリの例を示します。
select * from blink_test;
ワイドテーブルのマージと更新
このセクションでは、一般的なユースケースである、複数のストリームからのデータを単一の Hologres ワイドテーブルに書き込む方法について説明します。
列 A、B、C、D、E を持つ WIDE_TABLE という名前の Hologres ワイドテーブルがあり、列 A がプライマリキーであると仮定します。Flink では、1 つのストリームに列 A、B、C のデータが含まれ、別のストリームに列 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}パラメーターを実際の値に置き換えてください。 -
専用モードの Blink のバージョン 3.7 以降のみが、パーティションの自動作成をサポートしています。以前のバージョンを使用している場合は、事前に Hologres で子パーティションテーブルを作成する必要があります。そうしないと、データのインポートに失敗します。
-
-
専用モードの Blink でジョブを作成します。
以下に、専用モードの Blink でジョブを作成するためのステートメントの例を示します。
説明次の例は、Blink 専用モード 3.7 以降に適用されます。Blink 専用モードの 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 間のデータ型マッピングについては、「データ型の概要」をご参照ください。