Tair (Redis OSS-compatible) コネクタを使用すると、Flink SQL ストリーミングジョブで Tair をディメンションテーブルとして読み取ったり、結果テーブルとして書き込んだりできます。
コネクタの概要
| 項目 | 詳細 |
|---|---|
| テーブルタイプ | ディメンションテーブル、結果テーブル |
| 実行モード | ストリーミング |
| データフォーマット | STRING |
| メトリクス | ディメンションテーブル: なし。 シンクテーブル: numBytesOut、numRecordsOutPerSecond、numBytesOutPerSecond、currentSendTime。詳細については、「モニタリングメトリック」をご参照ください。 |
| API タイプ | SQL API |
| 結果テーブルでのデータの更新と削除 | サポートされています |
前提条件
開始する前に、以下が準備できていることを確認してください。
-
Tair (Redis OSS-compatible) インスタンス。「ステップ 1: インスタンスの作成」をご参照ください。
-
インスタンス用に設定されたホワイトリスト。「ステップ 2: ホワイトリストの設定」をご参照ください。
制限事項
-
配信セマンティクス:コネクタはベストエフォート型配信のみをサポートします。Exactly-once セマンティクスはサポートされていません。書き込み操作がべき等であることを確認してください。
-
ディメンションテーブルのデータ型:ディメンションテーブルは STRING および HASHMAP データのみを読み取ることができます。すべてのフィールドは STRING 型である必要があります。
-
ディメンションテーブルのプライマリキー:各ディメンションテーブルには、プライマリキーが 1 つだけ必要です。ディメンションテーブルの JOIN の ON 句では、プライマリキーに対する等価条件を使用する必要があります。
既知の問題
VVR 8.0.9 — バッファードライターのキャッシュのバグ:VVR 8.0.9 にはバッファードライターのキャッシュの問題が存在します。これを回避するには、結果テーブルの WITH 句で sink.buffer-flush.max-rows を 0 に設定します。
構文
CREATE TABLE redis_table (
col1 STRING,
col2 STRING,
PRIMARY KEY (col1) NOT ENFORCED -- 必須。
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'mode' = 'STRING' -- 結果テーブルに必須。
);
コネクタオプション
一般的なオプション
これらのオプションは、結果テーブルとディメンションテーブルの両方に適用されます。
| オプション | データ型 | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
connector |
STRING | はい | — | redis に設定します。 |
host |
STRING | はい | — | ApsaraDB for Redis データベースへの接続に使用される IP アドレス。可能な限り内部エンドポイントを使用してください。インターネット接続では、レイテンシーの増加や帯域幅の制限が発生する可能性があります。 |
port |
INT | いいえ | 6379 |
ポート番号。 |
password |
STRING | いいえ | (空の文字列) | アクセスパスワード。 |
dbNum |
INT | いいえ | 0 |
データベースシーケンス番号。 |
clusterMode |
BOOLEAN | いいえ | false |
データベースがクラスターモードであるかどうか。 |
hostAndPorts |
STRING | いいえ | — | "host1:port1,host2:port2" 形式のホストとポートのペア。clusterMode が true で、自己管理 Redis クラスターへの Jedis 接続に高可用性 (HA) が必要な場合に必須です。host と port よりも優先されます。clusterMode が true で HA が不要な場合は、host と port のみを設定して単一ノードを指定できます。 |
key-prefix |
STRING | いいえ | — | ディメンションテーブルからの読み取り時、または結果テーブルへの書き込み時にプライマリキーの値に追加されるプレフィックス。key-prefix-delimiter で指定されたデリミタが、プレフィックスとプライマリキーの値を区切ります。VVR 8.0.7 以降が必要です。 |
key-prefix-delimiter |
STRING | いいえ | — | キープレフィックスとプライマリキーの値の間のデリミタ。 |
connection.pool.max-total |
INT | いいえ | 8 |
接続プールによって割り当て可能な接続の最大数。VVR 8.0.9 以降が必要です。 |
connection.pool.max-idle |
INT | いいえ | 8 |
接続プール内のアイドル接続の最大数。 |
connection.pool.min-idle |
INT | いいえ | 0 |
接続プール内のアイドル接続の最小数。 |
connection.pool.lifo |
Boolean | いいえ | true |
アイドル接続が接続プールから LIFO 順で割り当てられるかどうか。有効な値:
注:このオプションは、Realtime Compute for Apache Flink エンジン VVR 11.8.0 以降でのみサポートされます。 |
connect.timeout |
DURATION | いいえ | 3000ms |
接続設定のタイムアウト。 |
socket.timeout |
DURATION | いいえ | 3000ms |
Redis サーバーからのデータ受信のタイムアウト。 |
cacert.filepath |
STRING | いいえ | — | SSL/TLS 証明書へのフルパス。ファイルは JKS 形式である必要があります。設定しない場合、SSL/TLS 暗号化は無効になります。暗号化を有効にするには、CA 証明書をダウンロードし、追加の依存関係としてアップロードします。これは /flink/usrlib ディレクトリに保存されます。例:'cacert.filepath' = '/flink/usrlib/ca.jks'。VVR 11.1 以降が必要です。 |
シンクテーブルのオプション
| オプション | データ型 | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
mode |
STRING | はい | — | 結果テーブルの Redis データ構造。STRING、LIST、SET、HASHMAP、SORTEDSET の 5 つの構造がサポートされています。DDL 文は選択した構造と一致する必要があります。「結果テーブルのデータ構造」をご参照ください。 |
flattenHash |
BOOLEAN | いいえ | false |
HASHMAP データを複数値モードで書き込むかどうか。true の場合、複数の非プライマリキーフィールドを宣言します。プライマリキーは Redis キーにマッピングされ、各非プライマリキーフィールド名はハッシュフィールドに、各フィールド値はハッシュ値にマッピングされます。false (単一値モード) の場合、3 つのフィールドを正確に宣言します。プライマリキーはキーに、最初の非プライマリキーフィールドはハッシュフィールドに、2 番目のフィールドはハッシュ値にマッピングされます。mode が HASHMAP の場合にのみ有効です。VVR 8.0.7 以降が必要です。 |
ignoreDelete |
BOOLEAN | いいえ | false |
リトラクションメッセージを無視するかどうか。true の場合、リトラクションメッセージは破棄されます。false の場合、リトラクションメッセージが受信されると、キーとそのデータが削除されます。 |
expiration |
LONG | いいえ | 0 |
挿入されたキーの Time-to-live (TTL) (ミリ秒単位)。0 は TTL を無効にします。 |
sink.buffer-flush.max-rows |
INT | いいえ | 200 |
バッファがフラッシュされる前に保持されるレコード (追加、変更、削除イベント) の最大数。clusterMode = false の場合は VVR 8.0.9 以降、clusterMode = true の場合は VVR 11.4.0 以降が必要です。 |
sink.buffer-flush.interval |
DURATION | いいえ | 1000ms |
バッファが非同期でフラッシュされる間隔。clusterMode = false の場合は VVR 8.0.9 以降、clusterMode = true の場合は VVR 11.4.0 以降が必要です。 |
ディメンションテーブルのオプション
| オプション | データ型 | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
mode |
STRING | いいえ | STRING |
ディメンションテーブルから読み取るデータ型。STRING は STRING データを読み取ります。HASHMAP はネストされたハッシュデータ (Key → Map\<Field, Value\>) を読み取ります。複数の非プライマリキー列を宣言し、プライマリキーは Redis キーに、各非プライマリキー列名はハッシュフィールドに、その値はフィールド値にマッピングされます。VVR 8.0.7 以降が必要です。HASHMAP データを単一値モードで読み取るには、代わりに hashName を設定します。 |
hashName |
STRING | いいえ | — | HASHMAP データを単一値モードで読み取る際に使用される固定ハッシュキー。設定した場合、2 つのフィールドを宣言します。プライマリキーはハッシュフィールドに、非プライマリキーはハッシュ値にマッピングされます。 |
cache |
STRING | いいえ | None |
キャッシュポリシー。None はキャッシュを無効にします。LRU はデータの一部をキャッシュします。キャッシュミスの場合、コネクタはディメンションテーブルをクエリします。ALL はデプロイメントが実行される前にディメンションテーブル全体をキャッシュにロードします。その後のすべてのルックアップはキャッシュを使用し、エントリの有効期限が切れるとキャッシュが再読み込みされます。「cache オプションの注意事項」をご参照ください。 |
cacheSize |
LONG | いいえ | 10000 |
キャッシュする行の最大数。cache が LRU の場合に必須です。 |
cacheTTLMs |
LONG | いいえ | — | キャッシュのタイムアウト (ミリ秒単位)。LRU の場合、エントリごとの有効期限を設定します (デフォルトでは有効期限なし)。ALL の場合、キャッシュの再読み込み間隔を設定します (デフォルトでは再読み込みなし)。cache が None の場合は効果がありません。 |
cacheEmpty |
BOOLEAN | いいえ | true |
空の結果 (一致なし) をキャッシュするかどうか。 |
cacheReloadTimeBlackList |
STRING | いいえ | — | ALL キャッシュポリシーが再読み込みしない期間。高トラフィックのイベント中に役立ちます。フォーマット:開始時刻と終了時刻を区切るには -> を、複数の期間を区切るには , を使用します。例:単一日 2017-10-24 14:00 -> 2017-10-24 15:00、日をまたぐ場合 2017-11-10 23:30 -> 2017-11-11 08:00、毎日繰り返す場合 12:00 -> 14:00, 22:00 -> 2:00 (VVR 11.1 以降が必要)。 |
async |
BOOLEAN | いいえ | false |
非同期ルックアップを有効にするかどうか。true の場合、結果は順不同で返されます。 |
cache オプションの注意事項
-
ALLは VVR 8.0.3 以降が必要です。 -
VVR 8.0.3 から VVR 11.1 (除く) まで、
cache = ALLは HASHMAP を単一値モードでのみ読み取ります。DDL の WITH 句で、hashNameをキー名に設定し、Field をプライマリキーとして、Value を非プライマリキー列として宣言します。 -
VVR 11.1 以降、
cache = ALLは HASHMAP の複数値モードをサポートします。Redis キーをプライマリキーとして指定し、各ハッシュフィールドに対して複数の非プライマリキー列を宣言します。WITH 句でmode = HASHMAPを設定します。 -
cacheオプションはcacheSizeおよびcacheTTLMsと一緒に使用する必要があります。
結果テーブルのデータ構造
各 Redis データ構造は、特定の DDL スキーマと書き込みコマンドに対応しています。
| データ構造 | DDL スキーマ | 書き込みコマンド |
|---|---|---|
| STRING | 2 つの列:key (STRING)、value (STRING) | set key value |
| LIST | 2 つの列:key (STRING)、value (STRING) | lpush key value |
| SET | 2 つの列:key (STRING)、value (STRING) | sadd key value |
| HASHMAP (単一値モード、デフォルト) | 3 つの列:key (STRING)、field (STRING)、value (STRING) | hmset key field value |
HASHMAP (複数値モード、flattenHash = true) |
複数の列:key (STRING)、次にハッシュフィールドごとに 1 つの列 — 各列名がフィールド名、その値がフィールド値 | hmset key col1 value1 col2 value2 ... |
| SORTEDSET | 3 つの列:key (STRING)、score (DOUBLE)、value (STRING) | zadd key score value |
ignoreDeleteオプションは、リトラクションメッセージの処理方法を制御します。trueに設定すると、削除操作はスキップされます。
データ型のマッピング
| スコープ | Tair (ApsaraDB for Redis) の型 | Flink の型 |
|---|---|---|
| すべてのテーブルタイプ | STRING | STRING |
| シンクテーブルのみ | SCORE | DOUBLE |
SCORE 型は SORTEDSET データと共に使用されます。ソートセット内の各値には DOUBLE のスコアが必要で、値はスコアの昇順でソートされます。
例
結果テーブルの例
すべての結果テーブルの例では、Kafka ソースから読み取り、Tair の結果テーブルに書き込みます。
STRING データの書き込み
この例では、user_id を Redis キーとして、login_time を Redis の値として使用します。
CREATE TEMPORARY TABLE kafka_source (
user_id STRING, -- ユーザー ID
login_time STRING -- ログイン時刻 (UNIX タイムスタンプ)
) WITH (
'connector' = 'kafka',
'topic' = 'user_logins',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
user_id STRING, -- Redis キー
login_time STRING, -- Redis の値
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'STRING',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
HASHMAP データの複数値モードでの書き込み
この例では、order_id を Redis キーとして使用し、product_name、quantity、amount を別々のハッシュフィールドとして書き込みます。
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- 注文 ID
product_name STRING, -- 製品名
quantity STRING, -- 製品数量
amount STRING -- 注文金額
) WITH (
'connector' = 'kafka',
'topic' = 'orders_topic',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
order_id STRING, -- Redis キー
product_name STRING, -- ハッシュフィールド:product_name
quantity STRING, -- ハッシュフィールド:quantity
amount STRING, -- ハッシュフィールド:amount
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'HASHMAP',
'flattenHash' = 'true',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
HASHMAP データの単一値モードでの書き込み
この例では、order_id を Redis キーとして、product_name をハッシュフィールドとして、quantity をハッシュ値として使用します。
CREATE TEMPORARY TABLE kafka_source (
order_id STRING, -- 注文 ID
product_name STRING, -- 製品名
quantity STRING -- 製品数量
) WITH (
'connector' = 'kafka',
'topic' = 'orders_topic',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_sink (
order_id STRING, -- Redis キー
product_name STRING, -- Redis ハッシュフィールド
quantity STRING, -- Redis ハッシュ値
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'mode' = 'HASHMAP',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>'
);
INSERT INTO redis_sink
SELECT * FROM kafka_source;
ディメンションテーブルの例
すべてのディメンションテーブルの例では、Tair ディメンションテーブルからユーザー情報をルックアップし、それを Kafka ストリームと結合します。
STRING データの読み取り
この例では、user_id を Redis キーとして使用し、user_name を Redis の値として取得します。
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- Redis キー
user_name STRING, -- Redis の値
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'mode' = 'STRING'
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
redis_user_id STRING,
user_name STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_id,
t2.user_name
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;
HASHMAP データの複数値モードでの読み取り
この例では、user_id を Redis キーとして使用し、複数のハッシュフィールド — user_name、email、register_time — を 1 回のルックアップで取得します。
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
click_time TIMESTAMP(3),
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- Redis キー
user_name STRING, -- ハッシュフィールド:user_name
email STRING, -- ハッシュフィールド:email
register_time STRING, -- ハッシュフィールド:register_time
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'mode' = 'HASHMAP'
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
user_name STRING,
email STRING,
register_time STRING,
click_time TIMESTAMP(3)
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_name,
t2.email,
t2.register_time,
t1.click_time
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;
HASHMAP データの単一値モードでの読み取り
この例では、hashName を介して設定された固定ハッシュキー (testkey) を使用します。user_id 列はハッシュフィールドに、user_name はハッシュ値にマッピングされます。
CREATE TEMPORARY TABLE kafka_source (
user_id STRING,
proctime AS PROCTIME()
) WITH (
'connector' = 'kafka',
'topic' = 'user_clicks',
'properties.bootstrap.servers' = '<yourKafkaBroker>',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TEMPORARY TABLE redis_dim (
user_id STRING, -- ハッシュフィールド
user_name STRING, -- ハッシュ値
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'redis',
'host' = '<yourHost>',
'port' = '<yourPort>',
'password' = '<yourPassword>',
'hashName' = 'testkey' -- 固定ハッシュキー
);
CREATE TEMPORARY TABLE blackhole_sink (
user_id STRING,
redis_user_id STRING,
user_name STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT
t1.user_id,
t2.user_id,
t2.user_name
FROM kafka_source AS t1
JOIN redis_dim FOR SYSTEM_TIME AS OF t1.proctime AS t2
ON t1.user_id = t2.user_id;