このトピックでは、ML_PREDICT 関数を使用して Flink で AI モデルを呼び出す方法について説明します。構文とパラメーター、呼び出しごとの設定、コンテンツタイプの設定、列レベルのパラメーター、およびテキスト、画像、マルチモーダルの推論の例を解説します。
クイックスタート
前提条件
-
Flink ワークスペースが作成されていること。詳細については、「Realtime Compute for Apache Flinkの有効化」をご参照ください。
-
Flink AI Service が有効化されていること。詳細については、「Flink AI Service (組み込みモデル)」をご参照ください。
-
Flink AI Service (組み込みモデル) には、VVR エンジンバージョン 11.7 以降が必要です。
次の例では、ML_PREDICT を使用して Flink の組み込みモデルを呼び出す方法を示します。 に移動し、ジョブを作成してコードを貼り付け、[Debug] をクリックします。
CREATE TEMPORARY TABLE text_source (
user_input STRING
) WITH ('connector' = 'datagen');
CREATE TEMPORARY TABLE result_sink (
user_input STRING,
ai_analysis STRING
) WITH ('connector' = 'print');
CREATE TEMPORARY MODEL text_model
INPUT (user_input STRING)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.6-flash',
'task' = 'chat/completions',
'system-prompt' = 'Rate the gibberish level of the input on a scale of 0 to 100'
);
INSERT INTO result_sink
SELECT user_input, content as ai_analysis FROM
ML_PREDICT(
TABLE text_source,
MODEL text_model,
DESCRIPTOR(user_input)
);
使用制限
-
Realtime Compute for Apache Flink エンジン VVR 11.1 以降が必要です。
-
一部のパラメーターは、Flink AI Service (組み込みモデル) を使用する場合にのみサポートされ、VVR 11.8.preview.2 以降が必要です。
-
ML_PREDICT 演算子のスループットは、Model Studio によるレート制限の対象となります。レート制限に達すると、ML_PREDICT 演算子でバックプレッシャーが発生し、タイムアウトエラーやジョブの再起動が発生する可能性があります。詳細については、Model Studio の「レート制限」をご参照ください。
-
content-types内のタイプ数と DESCRIPTOR の列数は、CREATE MODEL で定義された INPUT 列数と一致する必要があります。 -
型が
image_urlの列は、STRING である必要があります。型がmulti_image_urlsの列は、ARRAY<STRING>である必要があります。 -
base64 画像には、
data:image/<format>;base64,プレフィックスを含める必要があります。生の base64 文字列およびローカルファイルパスはサポートされていません。
構文
ML_PREDICT(TABLE <table_name>, MODEL <model_name>, DESCRIPTOR(<input_columns>) [, CONFIG => MAP[...]])
パラメーター
|
パラメーター |
データ型 |
説明 |
|
TABLE |
TABLE |
モデル推論に使用する入力データストリーム。物理テーブルまたはビューを指定できます。 |
|
MODEL |
MODEL |
登録済みモデルの名前。詳細については、「モデル設定」をご参照ください。 |
|
DESCRIPTOR() |
— |
モデル推論に使用する入力列。 説明
VVR 11.8.preview.2 以降では、複数の入力列がサポートされています。これは、Flink AI Service (組み込みモデル) を使用する場合にのみ利用可能です。DESCRIPTOR 列の数は、CREATE MODEL の INPUT 列の数と一致する必要があります。 |
|
CONFIG => MAP[...] |
MAP |
オプション。詳細については、「呼び出しごとの設定」をご参照ください。 説明
VVR 11.8.preview.2 以降で Flink AI Service (組み込みモデル) を使用する場合にのみサポートされます。 |
呼び出しごとの設定
ML_PREDICT を呼び出すときに、呼び出しごとの設定を指定できます。パラメーターがすでに CREATE MODEL で設定されている場合、呼び出しごとの値が優先されますが、MODEL の定義には永続化されません。
|
パラメーター |
説明 |
例 |
|
|
ユーザープロンプトを指定します。空の文字列を渡して、MODEL レベルの値をバイパスします。 |
|
|
|
単一列入力用のコンテンツタイプを指定します。 |
|
|
|
複数列の入力のコンテンツタイプを指定します。 |
|
|
|
列レベルのパラメーターを指定します。 |
|
|
|
追加のパラメーターを JSON 文字列として指定します。 |
|
コンテンツタイプのパラメーター
-
単一列の入力では、コンテンツタイプを指定するには
content-typeを使用します。サポートされている値はtextとimage_urlです。 -
複数カラム入力では、
content-typesを使用して各カラムのコンテンツタイプを指定します。サポートされている値:text、image_url、multi_image_urls。 -
content-typeとcontent-typesのうち1つのみ指定できます。使用できるオプションは、CREATE MODEL で設定された内容によって異なります。
|
CREATE MODEL の設定 |
許可される呼び出しごとのオプション |
許可されない呼び出しごとのオプション |
説明 |
|
|
|
|
|
|
|
|
|
たとえば、 |
|
|
|
— |
上記と同じ |
列レベルのパラメーター
|
パラメーター |
説明 |
値 |
例 |
|
|
入力画像またはビデオフレームの最小ピクセルしきい値を設定します。ピクセル数が |
|
|
|
|
入力画像またはビデオフレームの最大ピクセルしきい値を設定します。 |
|
|
|
|
動画から抽出された全フレームの合計ピクセル数 (1 フレームあたりのピクセル数 x 総フレーム数) を制限します。動画がこの制限を超えた場合、フレームは、各フレームを |
|
|
|
|
明示的なキャッシュを有効にします。 |
|
|
例
テキスト
次の例では、Flink AI Service の組み込みモデルを登録して使用し、入力テキストの感情を分類します。例 1 では、MODEL レベルのパラメーターを使用します。例 2 では、呼び出し時に user-prompt をオーバーライドします。
-- 組み込みモデルを登録します。content-type が指定されていない場合、デフォルトは text です
CREATE MODEL sentiment_model
INPUT (prompt STRING)
OUTPUT (response STRING)
WITH (
'provider' = 'openai-compat',
'task' = 'chat/completions',
'model' = 'qwen3.6-flash',
'system-prompt' = 'You are a sentiment classifier. Output one label: negative, positive, or neutral'
);
-- ソーステーブルを作成します:製品レビューのシミュレーションデータ
CREATE TEMPORARY VIEW input_table(id, content)
AS VALUES
(1, 'Great quality, soft fabric, fits perfectly'),
(2, 'Had loose threads on arrival, faded badly after one wash'),
(3, 'Received the item, looks as pictured'),
(4, 'Started pilling after two weeks, customer service refused returns'),
(5, 'Flattering fit, color is even better than the photo, already ordered a third one');
-- 結果テーブルを作成します
CREATE TEMPORARY TABLE output_table (
id INT,
content STRING,
sentiment STRING
) WITH (
'connector' = 'print'
);
-- ML_PREDICT を使用してリアルタイム推論を実行します
-- 例 1: MODEL で定義したパラメーターを使用します
INSERT INTO output_table
SELECT
id,
content,
response AS sentiment
FROM ML_PREDICT(
TABLE input_table,
MODEL sentiment_model,
DESCRIPTOR(content));
-- 例 2: 呼び出しごとのパラメーターを指定します
INSERT INTO output_table
SELECT
id,
content,
response AS sentiment
FROM ML_PREDICT(
TABLE input_table,
MODEL sentiment_model,
DESCRIPTOR(content),
MAP['user-prompt', 'Reply in Chinese']);
例 1 の出力:
|
id |
content |
sentiment |
1 |
Great quality, soft fabric, fits perfectly |
positive |
2 |
Had loose threads on arrival, faded badly after one wash |
negative |
3 |
Received the item, looks as pictured |
positive |
4 |
Started pilling after two weeks, customer service refused returns |
negative |
5 |
Flattering fit, color is even better than the photo, already ordered a third one |
positive |
例 2 の出力:
id |
content |
sentiment |
1 |
Great quality, soft fabric, fits perfectly |
正面 |
2 |
Had loose threads on arrival, faded badly after one wash |
負面 |
3 |
Received the item, looks as pictured |
正面 |
4 |
Started pilling after two weeks, customer service refused returns |
負面 |
5 |
Flattering fit, color is even better than the photo, already ordered a third one |
正面 |
画像
次の例では、画像を入力として受け取るモデルを登録し、入力画像を分類します。
-- content-type を 'image_url' に設定して組み込みモデルを登録します
CREATE TEMPORARY MODEL image_model
INPUT (prompt STRING)
OUTPUT (response STRING)
WITH (
'provider' = 'openai-compat',
'task' = 'chat/completions',
'model' = 'qwen3.6-flash',
'content-type' = 'image_url'
);
-- ML_PREDICT を使用してリアルタイム推論を実行します
INSERT INTO output_table
SELECT
id,
content,
response AS sentiment
FROM ML_PREDICT(
TABLE input_table,
MODEL image_model,
DESCRIPTOR(content));
テキスト + 単一画像
次の例では、テキストと画像の入力を使い、チャットベースの推論を行うマルチモーダルモデルを登録します。
CREATE MODEL vl_model
INPUT (text_input STRING, image_input STRING)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.5-plus',
'task' = 'chat/completions',
'content-types' = 'text;image_url'
);
INSERT INTO result_sink
SELECT content FROM TABLE(ML_PREDICT(
TABLE image_source,
MODEL vl_model,
DESCRIPTOR(text_input, image_input)
));
テキスト + 複数画像 (複数列)
各画像を個別の INPUT 列で渡します。image_url は、content-types 内の各画像列に指定します。
CREATE MODEL vl_model_multi
INPUT (prompt STRING, img1 STRING, img2 STRING)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.5-plus',
'task' = 'chat/completions',
'content-types' = 'text;image_url;image_url'
);
SELECT content FROM TABLE(ML_PREDICT(
TABLE my_source,
MODEL vl_model_multi,
DESCRIPTOR(prompt, img1, img2)
));
テキスト + 複数画像 (配列)
複数の画像を ARRAY<STRING> 列で渡します。content-types には multi_image_urls を使用します。
CREATE MODEL vl_model_array
INPUT (prompt STRING, images ARRAY<STRING>)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.5-plus',
'task' = 'chat/completions',
'content-types' = 'text;multi_image_urls'
);
SELECT content FROM TABLE(ML_PREDICT(
TABLE my_source,
MODEL vl_model_array,
DESCRIPTOR(prompt, images)
));
呼び出しごとのパラメーターの例
例 1:コンテンツタイプのパラメーターを指定する
-- CREATE MODEL で content-type / content-types が指定されていない場合、デフォルトは content-type = text です
CREATE MODEL model_single
INPUT (input STRING)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.5-plus',
'task' = 'chat/completions'
);
-- 呼び出し 1: MODEL で定義したパラメーターを使用します。コンテンツタイプは text です
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_text,
MODEL model_single,
DESCRIPTOR(input)
));
-- 呼び出し 2: 呼び出し時にコンテンツタイプを画像に上書きします
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_img,
MODEL model_single,
DESCRIPTOR(input),
MAP['content-type', 'image_url']
));
例 2:user-prompt と列レベルのパラメーターを指定する
-- CREATE MODEL でデフォルト設定を行います (複数列モデル)
CREATE MODEL vl_model
INPUT (text_input STRING, image_input STRING)
OUTPUT (content STRING)
WITH (
'provider' = 'openai-compat',
'model' = 'qwen3.5-plus',
'task' = 'chat/completions',
'content-types' = 'text;image_url',
'user-prompt' = 'Describe the image'
);
-- 呼び出し 1: デフォルト設定を使用します
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_a, MODEL vl_model, DESCRIPTOR(text_input, image_input)
));
-- 呼び出し 2: user-prompt と列レベルのパラメーターを上書きします
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_b, MODEL vl_model, DESCRIPTOR(text_input, image_input),
MAP[
'user-prompt', 'Answer in English',
'image_input.min_pixels', '100',
'image_input.max_pixels', '5000'
]
));
-- 呼び出し 3: content-types を上書きしてタイプの組み合わせを変更します (両方の列をテキストとして扱います)
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_c, MODEL vl_model, DESCRIPTOR(text_input, image_input),
MAP['content-types', 'text;text']
));
-- 呼び出し 4: content-types を上書きしてタイプの組み合わせを変更します (逆の順序)
SELECT content FROM TABLE(ML_PREDICT(
TABLE source_c, MODEL vl_model, DESCRIPTOR(text_input, image_input),
MAP['content-types', 'image_url;text']
));