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

Realtime Compute for Apache Flink:ベクトル化結果キャッシュ

最終更新日:Sep 22, 2026

このドキュメントでは、AI_EMBED 関数に外部キャッシュテーブルを設定し、同一の入力テキストに対するベクトル化結果を再利用することで、モデル呼び出しを削減し、レイテンシーを低減し、コストを節約する方法について説明します。

背景情報

AI_EMBED 関数は、CACHE_TABLE パラメーターを使用して、ベクトル化結果の永続的なキャッシュとして外部テーブルを指定します。実行時の動作は次のとおりです:

  • キャッシュヒット:非同期ルックアップで一致が見つかった場合、関数はモデル呼び出しをスキップし、結果を直接返します。

  • キャッシュミス:非同期ルックアップで一致が見つからなかった場合、関数はモデルを呼び出してベクトルを計算し、結果を返した後、非同期でキャッシュテーブルにライトバックします。

このメカニズムにより、同一の入力テキストを持つリクエストが繰り返された場合の冗長なモデル呼び出しを防ぎます。

制限事項

  • この機能は、Realtime Compute for Apache Flink VVR 11.8.preview1 以降のバージョンでのみサポートされています。

  • Stream Storage Apache Fluss Edition の有効化。キャッシュテーブルは Fluss プライマリキーテーブルである必要があります。

キャッシュテーブルの作成

キャッシュテーブルのプライマリキーの型は、CACHE_KEY の型と互換性を持つ必要があります。

明示的なキー指定モード

プライマリキーの型は、CACHE_KEY で指定された列の型と互換性を持つ必要があります。たとえば、ソーステーブルの cache_key が INT の場合:

CREATE TABLE `fluss`.`default`.`ai_cache` (
  cache_key   INT NOT NULL,
  embedding   ARRAY<FLOAT>,PRIMARY KEY (cache_key) NOT ENFORCED
);

自動生成キーモード

プライマリキーは、SHA256 ハッシュ値を格納するために STRING 型である必要があります。

CREATE TABLE `fluss`.`default`.`ai_cache_sha` (
  cache_key   STRING NOT NULL,
  embedding   ARRAY<FLOAT>,PRIMARY KEY (cache_key) NOT ENFORCED
);

構文

AI_EMBED(
  MODEL       => MODEL <MODEL NAME>,
  INPUT       => <INPUT COLUMN NAME>
  [, DIMENSION   => <DIMENSION VALUE>]
  [, CACHE_TABLE => TABLE <catalog>.<db>.<table>
     [, CACHE_KEY => DESCRIPTOR(<column>)]]
  [, CONFIG      => MAP[...]]
)

パラメーター

パラメーター

必須

説明

MODEL

はい

モデルへの参照です。詳細については、「モデルの登録」をご参照ください。

INPUT

はい

入力テキスト列です。

CACHE_TABLE

いいえ

キャッシュテーブルへの参照です。構文は TABLE <catalog>.<db>.<table> です。

CACHE_KEY

いいえ

キャッシュキー列です。構文は DESCRIPTOR(<column>) です。指定しない場合、システムは SHA256 ハッシュを使用して複合キーを自動的に生成します。

CONFIG

いいえ

ランタイム設定です。

CACHE_KEY モード

明示的な指定 (推奨):入力テーブルの特定の列の値をルックアップキーとして使用します。このモードは、データに一意のキーが既にある場合に推奨します。

CACHE_KEY => DESCRIPTOR(cache_key)

自動生成:CACHE_KEY を渡さない場合、システムは次のルールに基づいてキーを生成します。

-- CACHE_KEY が指定されていない場合、システムはキーを自動的に生成します:
-- key = LOWER(SHA256(model_id || ':' || dimension || ':' || input_text))
CACHE_TABLE => TABLE `fluss`.`default`.`ai_cache1`

このモードでは、実行時に SHA256 ハッシュを計算するオーバーヘッドが発生しますが、モデル、ディメンション、および入力テキストの各組み合わせに対して一意のキーが自動的に保証されます。

パフォーマンスチューニング

チューニングパラメーター

SET 'table.exec.async-ml-predict.max-concurrent-operations' = '100';
SET 'table.exec.async-ml-predict.output-mode' = 'ALLOW_UNORDERED';
SET 'table.exec.async-ml-predict.timeout' = '120s';

パラメーター

パラメーター

デフォルト値

説明

max-concurrent-operations

10

キャッシュルックアップまたはモデル呼び出しを並行して実行するレコードの最大数です。この値を大きくすると、スループットが向上する可能性があります。ただし、値を大きくしすぎると、モデルサービスに過負荷がかかったり、過剰なメモリ使用量が発生したりする可能性があります。

output-mode

ORDERED

出力順序モードです:

  • ORDERED:出力が入力と同じ順序であることを保証します。これは、結合やウィンドウ関数など、順序に依存する後続の操作に適しています。

  • ALLOW_UNORDERED:レコードが完了するとすぐに出力します。これは、insert-only のストリームに推奨します。

timeout

3 min

ルックアップと予測を合わせた最大許容時間です。この値が低すぎると、キャッシュミス時の通常操作がタイムアウトする可能性があります。

推奨事項:ALLOW_UNORDERED を使用すると、次の理由からパフォーマンスが大幅に向上するため、推奨されます:

  • キャッシュヒットは通常 1 回のルックアップしか必要としませんが、キャッシュミスは追加のモデル呼び出しを必要とします。

  • ORDERED モードでは、1 回のキャッシュミスが後続のすべてのキャッシュヒットの出力をブロックします。

  • ALLOW_UNORDERED モードでは、キャッシュヒットは先行するミスの完了を待たず、即座に出力されます。

説明

insert-only のストリームのみが ALLOW_UNORDERED モードを使用します。撤回 (retraction) を含むストリームの場合、ALLOW_UNORDERED が指定されていても、システムは ORDERED モードにフォールバックします。

モニタリングメトリクス

ジョブの実行中に、Flink メトリクスを使用してキャッシュの有効性をモニタリングできます。

メトリクス名

説明

注

ai_function_cache.cache_hit

キャッシュヒット数

ヒット率 = ヒット / (ヒット + ミス)。

ai_function_cache.cache_miss

キャッシュミス数

この値が一貫して高い場合は、キャッシュがウォームアップされていないか、キーが不安定であることを示します。

ai_function_cache.writeback_attempt

ライトバック試行回数

この値は、ミス数に近い値になるはずです。

ai_function_cache.writeback_skipped_null_output

null 出力のためにスキップされたライトバック数

この値が高い場合は、モデルに潜在的な問題があることを示します。

例

Fluss プライマリキーテーブルの例:

-- キャッシュテーブル (Fluss) を事前に作成します
CREATE TABLE `fluss`.`default`.`embedding_cache` (
  cache_key   STRING NOT NULL,
  embedding   ARRAY<FLOAT>,
  PRIMARY KEY (cache_key) NOT ENFORCED
);

この例では、Fluss プライマリキーテーブルをキャッシュとして使用し、出力モードを ALLOW_UNORDERED に設定してスループットを向上させます。

-- チューニングパラメーター
SET 'table.exec.async-ml-predict.max-concurrent-operations' = '100';
SET 'table.exec.async-ml-predict.output-mode' = 'ALLOW_UNORDERED';

-- モデルの作成
CREATE TEMPORARY MODEL embedding_model
  INPUT (text STRING)
  OUTPUT (embedding ARRAY<FLOAT>)
WITH (
  'provider' = 'openai-compat',
  'task' = 'embeddings',
  'model' = 'text-embedding-v4',
  'dimension' = '1024'
);

-- ビジネスロジックのクエリ
INSERT INTO result_sink
SELECT id, content, embedding
FROM source_table, LATERAL TABLE(AI_EMBED(
  MODEL       => MODEL embedding_model,
  INPUT       => content,
  CONFIG      => MAP['async', 'true'],
  CACHE_TABLE => TABLE `fluss`.`default`.`embedding_cache`
));