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

Realtime Compute for Apache Flink:Tair (Redis OSS-compatible) コネクタ

最終更新日:Jul 31, 2026

Tair (Redis OSS-compatible) コネクタを使用すると、Flink SQL ストリーミングジョブで Tair をディメンションテーブルとして読み取ったり、結果テーブルとして書き込んだりできます。

コネクタの概要

項目 詳細
テーブルタイプ ディメンションテーブル、結果テーブル
実行モード ストリーミング
データフォーマット STRING
メトリクス ディメンションテーブル: なし。 シンクテーブル: numBytesOutnumRecordsOutPerSecondnumBytesOutPerSecondcurrentSendTime。詳細については、「モニタリングメトリック」をご参照ください。
API タイプ SQL API
結果テーブルでのデータの更新と削除 サポートされています

前提条件

開始する前に、以下が準備できていることを確認してください。

制限事項

  • 配信セマンティクス:コネクタはベストエフォート型配信のみをサポートします。Exactly-once セマンティクスはサポートされていません。書き込み操作がべき等であることを確認してください。

  • ディメンションテーブルのデータ型:ディメンションテーブルは STRING および HASHMAP データのみを読み取ることができます。すべてのフィールドは STRING 型である必要があります。

  • ディメンションテーブルのプライマリキー:各ディメンションテーブルには、プライマリキーが 1 つだけ必要です。ディメンションテーブルの JOIN の ON 句では、プライマリキーに対する等価条件を使用する必要があります。

既知の問題

VVR 8.0.9 — バッファードライターのキャッシュのバグ:VVR 8.0.9 にはバッファードライターのキャッシュの問題が存在します。これを回避するには、結果テーブルの WITH 句で sink.buffer-flush.max-rows0 に設定します。

構文

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" 形式のホストとポートのペア。clusterModetrue で、自己管理 Redis クラスターへの Jedis 接続に高可用性 (HA) が必要な場合に必須です。hostport よりも優先されます。clusterModetrue で HA が不要な場合は、hostport のみを設定して単一ノードを指定できます。
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 順で割り当てられるかどうか。有効な値:
  • true:LIFO。プールは最後に返された接続を最初に割り当てます。
  • false:FIFO。最も最近使用されていないアイドル接続が最初に割り当てられます。
プロキシの背後にある Redis クラスターインスタンスの場合、これを false に設定すると、プロキシ間の負荷をより均等に分散させ、接続の偏りを回避できます。ただし、FIFO は接続をビジー状態に保つ傾向があり、総接続数を増加させる可能性があるため、接続数が多いシナリオでは慎重に評価してください。

注:このオプションは、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 番目のフィールドはハッシュ値にマッピングされます。modeHASHMAP の場合にのみ有効です。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 キャッシュする行の最大数。cacheLRU の場合に必須です。
cacheTTLMs LONG いいえ キャッシュのタイムアウト (ミリ秒単位)。LRU の場合、エントリごとの有効期限を設定します (デフォルトでは有効期限なし)。ALL の場合、キャッシュの再読み込み間隔を設定します (デフォルトでは再読み込みなし)。cacheNone の場合は効果がありません。
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_namequantityamount を別々のハッシュフィールドとして書き込みます。

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_nameemailregister_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;