Python DataFrame API は、Flink ジョブを作成するための DataFrame スタイルの Python インターフェイスを提供します。`filter`、`project`、`join`、`aggregate` などの基本的な API を提供し、関係代数スタイルでデータ処理ロジックを表現できます。また、DataFrame API では、標準の演算子では表現できないビジネスロジックをカバーするために、Python でユーザー定義関数を作成することもできます。さらに、DataFrame API は、データパイプライン内で大規模言語モデルに直接接続する AI/LLM 機能を提供し、分類、感情分析、抽出、翻訳、要約、埋め込みのためのすぐに使える AI 関数を備えています。このトピックでは、ジョブ開発のリファレンスとしてサポートされている API について説明します。完全な API ドキュメントについては、「PyFlink DataFrame」をご参照ください。
DataFrame 操作
次の表は、DataFrame のコア操作をまとめたものです。
|
カテゴリ |
API |
|
構築/作成 |
from_table、from_pandas、from_arrow、from_dict、from_records、range |
|
属性 |
|
|
射影/列操作 |
|
|
フィルター |
|
|
集約 |
|
|
結合/和集合 |
|
|
行マッピング |
|
|
展開 |
|
|
重複排除 |
|
|
集合演算 |
|
|
制限/ページネーション |
|
|
パーティション |
|
|
パイプ |
|
|
Null 処理 |
|
|
出力/収集 |
|
|
デバッグ/実行計画 |
|
|
SQL |
式ヘルパー関数
|
関数 |
説明 |
|
列参照式を作成します。 |
|
|
リテラル式を作成します。 |
データ型
DataType は DataFrame の列の型を記述するために使用されます。
|
型カテゴリ |
メソッド |
|
ブール値 |
|
|
整数 |
|
|
浮動小数点/固定小数点 |
|
|
文字列 |
|
|
バイナリ |
|
|
日付/時刻 |
DataType.date、DataType.time、DataType.timestamp、DataType.timestamp_ltz |
|
複合 |
|
|
特殊 |
|
|
マルチモーダル |
|
|
Null 許容性修飾子 |
DataType は、データソースのスキーマや UDF の出力型などを指定するために使用できます。例えば、DataType を使用して Kafka ソースのスキーマを指定します。
import pyflink.dataframe as pf
from pyflink.dataframe import DataType
df = pf.read_kafka(
"localhost:9092",
topic="user_events",
schema={
"user_id": DataType.string(),
"event_type": DataType.string(),
"amount": DataType.decimal(10, 2),
"tags": DataType.list(DataType.string()),
"event_time": DataType.timestamp_ltz(3),
},
format="json",
startup_mode="earliest-offset",
)
DataType.tensor は可変形状テンソルをサポートします。UDF では、組み込みメソッドを使用して可変形状テンソルと NumPy ndarray を相互に変換できます。
from pyflink.dataframe.tensor import variable_tensors_to_arrow
result_type = DataType.struct({
"id": DataType.int64(),
# 形状指定なし:可変長テンソル
"tensor": DataType.tensor("float32"),
})
def transpose_batch_tensor(batch):
# to_numpy() を呼び出して、可変長テンソル列をオブジェクト ndarray に変換します。
# 各要素は独立した ndarray であり、形状は異なる場合があります。
tensors = batch["tensor"].to_numpy()
transposed = [
None if tensor is None else tensor.T
for tensor in tensors
]
return {
"id": batch["id"],
# 可変長 ndarray 列を Arrow の可変長テンソル列に戻します
"tensor": variable_tensors_to_arrow(
transposed,
dtype="float32",
),
}
可変形状テンソルは VVR 11.9.preview.1 からサポートされています。
I/O
DataFrame は、一般的なデータソースの読み取りと書き込みのための豊富な I/O API セットを提供します。また、`read_generic` または `write_generic` 関数を使用して、サポートされている他のコネクタやカスタムコネクタと連携してデータの読み書きを行うこともできます。
データの読み取り
|
ソース |
API |
|
Parquet ファイル |
|
|
JSON ファイル |
|
|
Kafka |
|
|
MaxCompute (ODPS) |
|
|
Paimon |
|
|
SLS (Simple Log Service) |
|
|
Hologres |
|
|
Milvus |
|
|
ビデオフレーム |
|
|
汎用/カスタムコネクタ |
データの書き込み
|
シンク |
API |
|
Parquet ファイル |
|
|
JSON ファイル |
|
|
Kafka |
|
|
MaxCompute (ODPS) |
|
|
Paimon |
|
|
SLS (Simple Log Service) |
|
|
Hologres |
|
|
Milvus |
|
|
プリント |
VVR 11.9.preview.1 からサポートされています。 |
|
汎用/カスタムコネクタ |
単一のジョブで複数回のデータ書き込みが必要な場合は、create_statement_set を使用して StatementSet を作成し、すべての書き込み関数をまとめて実行します。
カタログの使用
DataFrame API はカタログを使用し、カタログテーブルからの読み取りや書き込みができます。
|
API |
説明 |
|
カタログの作成 |
|
|
登録済みのカタログオブジェクトの取得 |
|
|
現在のカタログの切り替え |
|
|
現在のカタログの取得 |
|
|
カタログの一覧表示 |
|
|
現在のデータベースの切り替え |
|
|
現在のデータベースの取得 |
|
|
現在のカタログ内のデータベースの一覧表示 |
|
|
カタログテーブルを DataFrame として読み取る |
|
|
DataFrame をカタログテーブルに書き込む |
ユーザー定義関数
ユーザー定義スカラー関数 udf
ユーザー定義スカラー関数は、各入力行を単一の結果値に変換します。以下の 4 種類があります。
-
同期行 UDF:一度に 1 行を処理する通常の Python 関数です。
-
非同期行 UDF:外部サービスの呼び出しや I/O 集中型のロジックに適しており、複数の行が同時に I/O を実行することでスループットを向上させます。
-
同期ベクトル UDF:データをバッチで受け取り、返すことで、Python と Flink ランタイム間の行ごとのオーバーヘッドを削減します。純粋な計算に適しています。
-
非同期ベクトル UDF:バッチ処理のスループットと非同期 I/O の同時実行性を組み合わせます。
スカラー関数の登録
udf デコレーターは、Python 関数を DataFrame で使用するスカラー関数として登録します。
-
API 名:udf。
udf( func=None, *, return_dtype=None, deterministic=True, name=None, func_type=None, concurrency=None, batch_size=None, num_gpus=None, gpu_type=None ) -
説明:Python 関数を DataFrame のスカラー UDF として登録します。通常の関数、非同期関数、ScalarFunction / AsyncScalarFunction のサブクラス、および呼び出し可能クラスをサポートします。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
func
Callable / Class
いいえ
ラップする Python 関数、ScalarFunction インスタンス/サブクラス、または呼び出し可能クラス。省略した場合は、デコレーターが返されます。
return_dtype
DataType / str / type
いいえ
戻り値の型。
DataTypeインスタンス (例:DataType.int64())、Python 型 (例:int)、または SQL 型文字列 (例:'BIGINT') を指定できます。省略した場合、関数の型ヒントから自動的に推論されます。deterministic
Boolean
いいえ
関数が決定論的かどうか。デフォルト:
True。name
String
いいえ
UDF の名前。デフォルトは関数名です。
func_type
String
いいえ
実行フォーマット。オプション:
"general"、"pandas"、"arrow"。省略した場合、関数の型ヒントから自動的に検出されます。concurrency
int
いいえ
UDF 演算子の並列度。異なる同時実行数値を持つ UDF は、別々の演算子に分割されます。
batch_size
int
いいえ
バッチあたりの最大要素数。バッチ UDF (pandas / arrow モード) にのみ適用されます。
num_gpus
float
いいえ
UDF にリクエストする GPU の数 (例:0.5、1)。設定すると、UDF は別の演算子で実行され、他の UDF (他の GPU UDF を含む) とは連鎖しません。
gpu_type
String
いいえ
GPU の種類。num_gpus が指定されている場合に必須です。
-
戻り値
with_column、with_columns、map、map_batchesなどの操作で使用できる DataFrameUDFWrapper オブジェクト。
同期行 UDF
<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_column.html#pyflink.dataframe.DataFrame.with_column" id="9bd1023ff8zhv">with_column</a> または <a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_columns.html#pyflink.dataframe.DataFrame.with_columns" id="f6f90f4fb81f2">with_columns</a> での使用
スカラー UDF は 1 つ以上の列を受け取り、単一の列を返します。Python の型ヒントを使用して戻り値の型を自動推論するか、return_dtype を使用して明示的に指定します。UDF を定義する際に、concurrency を使用して同時実行数を設定できます。
VVR 11.9.preview.1 から、呼び出し可能クラスを使用して UDF を宣言できます。UDF の初期化と実行ロジックを __init__ と __call__ で定義します。
from pyflink.dataframe import col, udf
orders = pf.from_records(
[
(1001, 1, 99.9, "PAID"),
(1002, 2, 35.5, "CREATED"),
(1003, 1, 188.0, "PAID"),
(1004, 3, 88.8, None),
],
schema=["order_id", "user_id", "amount", "status"],
)
@udf
def normalize_status(status: str) -> str:
if status is None:
return "UNKNOWN"
return status.strip().upper()
@udf(concurrency=32)
def order_tag(amount: float, status: str) -> str:
"""複数の列を入力として受け取ります"""
if status == "PAID" and amount is not None and amount >= 100:
return "high_value_paid"
return "normal"
@udf
class ModelPredictor:
def __init__(self):
# 初期化はタスクマネージャーで実行されます
self._model = load_model()
def __call__(self, amount: float) -> str:
return self._model.predict(amount)
# with_column は単一の列を追加します
orders_with_status = orders.with_column(
"status_norm", normalize_status(col("status"))
)
# with_columns は一度に複数の列を追加します
orders_with_flags = orders.with_columns(
status_norm=normalize_status(col("status")),
tag=order_tag(col("amount"), col("status")),
)
# ステートフル UDF を使用します
orders_with_predictions = orders.with_column(
"prediction", ModelPredictor(col("amount"))
)
<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.map.html#pyflink.dataframe.DataFrame.map" id="39c01bde1d7d5">map</a> での使用
map は完全な行を受け取り、新しい行を返します。入力列は関数内で名前によってアクセスされます。戻り値の型は、DataType.struct で明示的に宣言するか、TypedDict の戻り値の型ヒントを通じて自動推論で宣言できます。udf をデコレーターとして使用せずに、Python 関数を直接 map に渡すことができます。
return_dtype を使用して明示的に宣言
def build_order_feature(row):
amount = row["amount"] or 0.0
return {
"order_id": row["order_id"],
"user_id": row["user_id"],
"feature": f"{row['status']}:{'large' if amount >= 100 else 'normal'}",
}
order_features = orders.map(
build_order_feature,
return_dtype=DataType.struct({
"order_id": DataType.int64(),
"user_id": DataType.int64(),
"feature": DataType.string(),
}),
)
TypedDict を使用して自動検出
from typing import TypedDict
class OrderFeature(TypedDict):
order_id: int
user_id: int
feature: str
def build_order_feature_typed(row) -> OrderFeature:
amount = row["amount"] or 0.0
return {
"order_id": row["order_id"],
"user_id": row["user_id"],
"feature": f"{row['status']}:{'large' if amount >= 100 else 'normal'}",
}
order_features_typed = orders.map(build_order_feature_typed)
同期ベクトル UDF
<a data-init-id="9bd1023ff8zhv" href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_column.html#pyflink.dataframe.DataFrame.with_column" id="e6c50a17edh0o">with_column</a> または <a data-init-id="f6f90f4fb81f2" href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_columns.html#pyflink.dataframe.DataFrame.with_columns" id="8a50565ecf76z">with_columns</a> での使用
ベクトル化 UDF はデータをバッチで受け取り、返すことで、Python と Flink ランタイム間の行ごとのオーバーヘッドを削減します。純粋な計算に適しています。Pandas (pandas.Series) と Arrow (pyarrow.Array) フォーマットをサポートします。UDF のロジックに基づいて選択してください。batch_size を使用してバッチサイズを制御します。UDF によって返される行数は、入力行数と等しくなければなりません。
Pandas フォーマット
import pandas as pd
@udf(return_dtype=DataType.float64(), concurrency=32, batch_size=64)
def scale_amount_pandas(amounts: pd.Series) -> pd.Series:
return amounts * 100.0
orders_scaled_pandas = orders.with_column(
"amount_scaled", scale_amount_pandas(col("amount")),
)
Arrow フォーマット
import pyarrow as pa
import pyarrow.compute as pc
@udf(return_dtype=DataType.float64(), concurrency=32, batch_size=64)
def scale_amount_arrow(amounts: pa.Array) -> pa.Array:
return pc.multiply(amounts, 100.0)
orders_scaled_arrow = orders.with_column(
"amount_scaled", scale_amount_arrow(col("amount")),
)
<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.map_batches.html#pyflink.dataframe.DataFrame.map_batches" id="d500967377ztd">map_batches</a> での使用
map_batches は、行全体に対してベクトル化処理を実行します。batch_format="pandas" の場合、関数の入力と出力は dict[str, pandas.Series] です。batch_format="arrow" の場合、入力と出力は dict[str, pyarrow.Array] です。udf デコレーターでラップせずに、Python 関数を直接 map_batches に渡すことができます。
Pandas フォーマット
def score_batch_pandas(batch: dict[str, pd.Series]) -> dict[str, pd.Series]:
amount = batch["amount"].fillna(0.0)
return {
"order_id": batch["order_id"],
"score": (amount / 100.0).clip(0.0, 1.0),
}
order_scores_pandas = orders.map_batches(
score_batch_pandas,
batch_format="pandas",
batch_size=1024,
return_dtype=DataType.struct(
{
"order_id": DataType.int64(),
"score": DataType.float64(),
}
),
)
Arrow フォーマット
def score_batch_arrow(batch: dict[str, pa.Array]) -> dict[str, pa.Array]:
amount = pc.if_else(pc.is_null(batch["amount"]), 0.0, batch["amount"])
raw_score = pc.divide(pc.cast(amount, pa.float64()), 100.0)
score = pc.if_else(
pc.less(raw_score, 0.0),
0.0,
pc.if_else(pc.greater(raw_score, 1.0), 1.0, raw_score),
)
return {
"order_id": batch["order_id"],
"score": score,
}
order_scores_arrow = orders.map_batches(
score_batch_arrow,
batch_format="arrow",
batch_size=1024,
return_dtype=DataType.struct(
{
"order_id": DataType.int64(),
"score": DataType.float64(),
}
),
)
非同期行 UDF
非同期 UDF を使用すると、複数の行で I/O 操作を同時に実行でき、スループットが向上します。これらは、外部サービスを呼び出すロジックや、I/O 負荷が高いロジックに適しています。非同期の行ベース UDF は、<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_column.html#pyflink.dataframe.DataFrame.with_column" id="a4862ccb65o83">with_column</a>、<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.with_columns.html#pyflink.dataframe.DataFrame.with_columns" id="1b48944541kmy">with_columns</a>、または map で呼び出すことができます。
import asyncio
# `with_column` での使用
@udf(concurrency=32)
async def query_region(user_id: int) -> str:
await asyncio.sleep(0.01) # 非同期 I/O をシミュレート
return f"region_for_{user_id}"
orders_with_region = orders.with_column(
"region", query_region(col("user_id")),
)
# `map` での使用
class OrderWithRegion(TypedDict):
order_id: int
region: str
async def enrich_order_region(row) -> OrderWithRegion:
await asyncio.sleep(0.01) # 非同期 I/O をモック
return {
"order_id": row["order_id"],
"region": f"region_for_{row['user_id']}",
}
orders_with_region = orders.map(enrich_order_region)
非同期ベクトル UDF
非同期ベクトル化 UDF は、バッチ処理のスループットと非同期 I/O の同時実行性を組み合わせます。batch_size を使用して各バッチのサイズを制御します。非同期ベクトル化 UDF は、with_column、with_columns、または map_batches で呼び出すことができます。
# `with_column` での使用
@udf(return_dtype=DataType.string(), concurrency=32, batch_size=64)
async def batch_enrich(statuses: pd.Series) -> pd.Series:
async def enrich_one(s):
await asyncio.sleep(0.01) # 非同期 API 呼び出しをシミュレート
return f"enriched_{s}"
tasks = [enrich_one(s) for s in statuses]
results = await asyncio.gather(*tasks)
return pd.Series(results)
orders_enriched = orders.with_column(
"status_enriched", batch_enrich(col("status")),
)
# `map_batches` での使用
async def enrich_order_batch(
batch: dict[str, pd.Series],
) -> dict[str, pd.Series]:
async def query_region(user_id):
await asyncio.sleep(0.01) # 非同期 I/O をシミュレート
return f"region_for_{user_id}"
regions = await asyncio.gather(
*(query_region(user_id) for user_id in batch["user_id"])
)
return {
"order_id": batch["order_id"],
"region": pd.Series(
regions,
index=batch["user_id"].index,
),
}
orders_with_region = orders.map_batches(
enrich_order_batch,
batch_format="pandas",
batch_size=64,
concurrency=32,
return_dtype=DataType.struct({
"order_id": DataType.int64(),
"region": DataType.string(),
}),
)
map および map_batches での非同期 UDF の呼び出しは、VVR 11.9.preview.1 からサポートされています。
ユーザー定義のテーブル関数 udtf
ユーザー定義のテーブル関数は、1 つの入力行を任意の数の出力行に展開し、各出力行には 1 つ以上の結果列を含めることができます。
テーブル関数の登録
udtf デコレーターは、Python 関数を DataFrame で使用するテーブル関数として登録します。
-
API 名:udtf。
udtf( func=None, *, return_dtype=None, deterministic=True, name=None, concurrency=None, num_gpus=None, gpu_type=None ) -
説明:Python 関数、
TableFunctionサブクラスまたはインスタンス、または呼び出し可能クラスを DataFrame で使用するテーブル関数として登録します。 -
パラメーター
パラメーター名
パラメーターの型
必須
説明
func
Callable / Class
いいえ
ラップする Python 関数、TableFunction インスタンス/サブクラス、または呼び出し可能クラス。省略した場合は、デコレーターが返されます。
return_dtype
DataType / str / type
いいえ
各出力行の型。複数列の出力には、
DataType.struct({...})またはTypedDictを使用します。deterministic
Boolean
いいえ
関数が決定論的かどうか。デフォルト:
True。name
String
いいえ
UDTF の名前。デフォルトは関数名です。
concurrency
int
いいえ
UDTF 演算子の並列度。
num_gpus
float
いいえ
UDTF にリクエストする GPU の数 (例:0.5、1)。設定すると、UDTF は別の演算子で実行され、他の UDF (他の GPU UDF を含む) とは連鎖しません。
gpu_type
String
いいえ
GPU の種類。num_gpus が指定されている場合に必須です。
-
戻り値
flat_mapまたはjoin_lateralで使用できる DataFrameUDTFWrapper オブジェクト。
join_lateral での使用
join_lateral は、元の入力行の後に UDTF の結果を追加します。デフォルトでは、UDTF の出力がない入力行は削除されます。ignore_empty=False を設定すると、UDTF の出力がない入力行を保持できます。
from typing import Iterator, Tuple
from pyflink.dataframe import udtf
# 型ヒントと .alias() を使用して、戻り値の型と列名を指定します
@udtf
def split_words(text: str) -> Iterator[Tuple[str, int]]:
for word in text.split():
yield word, len(word)
words_with_source = texts.join_lateral(
split_words(col("text")).alias("word", "word_length"),
ignore_empty=False,
)
# または return_dtype を使用します
@udtf(return_dtype=DataType.struct({
"word": DataType.string(),
"word_length": DataType.int32(),
}))
def split_words_with_length(text):
for word in text.split():
yield {"word": word, "word_length": len(word)}
words_with_detail = texts.join_lateral(
split_words_with_length(col("text")),
)
flat_map での使用
<a href="https://ververica-flink-docs.readthedocs.io/en/latest/reference/pyflink.dataframe/api/pyflink.dataframe.DataFrame.flat_map.html#pyflink.dataframe.DataFrame.flat_map" id="35bfeae711bka">flat_map</a> は、各行で UDTF を呼び出します。結果には UDTF の出力列のみが含まれ、元の入力列は保持されません。Python 関数は、udtf デコレーターでラップすることなく、flat_map に直接渡すことができます。
from typing import Any, Dict, Iterator, TypedDict
# TypedDict を使用した型ヒントで、戻り値の型と列名を指定します
class Word(TypedDict):
word: str
length: int
def split_row(row: Dict[str, Any]) -> Iterator[Word]:
for word in row["text"].split():
yield {"word": word, "length": len(word)}
words = texts.flat_map(split_row)
# または return_dtype を使用します
def split_row_with_dtype(row):
for word in row["text"].split():
yield {"word": word}
words_with_dtype = texts.flat_map(
split_row_with_dtype,
return_dtype=DataType.struct({"word": DataType.string()}),
)
AI / LLM 関数
DataFrame は、df.llm アクセサーを通じて一連の組み込み AI 関数を提供します。詳細な API リファレンスについては、「AI/LLM」をご参照ください。
プロバイダーの構成
AI 関数を使用する前に、まずモデルプロバイダーを登録します。
|
API |
説明 |
|
モデルプロバイダーを登録します。 |
|
|
デフォルトのプロバイダーを設定します (複数プロバイダーのシナリオ向け)。 |
|
|
登録済みのプロバイダーを一覧表示します。 |
サポートされているプロバイダーの種類:
|
プロバイダー |
ユースケース |
|
OpenAI、DeepSeek、DashScope (Model Studio)、およびすべての OpenAI 互換インターフェイス。 |
|
|
マルチモーダル埋め込みをサポートする Alibaba Cloud DashScope。 |
|
|
NVIDIA Triton Inference Server |
|
|
登録済みのカスタムモデルプロバイダーを使用します。 |
set_model_provider
-
API 名:set_model_provider。
set_model_provider( name_or_provider, provider=None, **options ) -
説明:グローバルな Model Provider 構成を登録します。3 つの呼び出しパターンがサポートされています:
ModelProviderインスタンスを渡す、名前とModelProviderインスタンスを渡す、または名前とキーワード引数を渡す。登録された名前は一意でなければなりません。既存の名前でプロバイダーを登録するとValueErrorが発生します。 -
パラメーター
パラメーター名
パラメーターの型
必須
説明
name_or_provider
ModelProvider / String
はい
ModelProviderインスタンス (プロバイダー識別子で自動的に登録される) またはカスタム名の文字列。provider
ModelProvider
いいえ
ModelProviderインスタンス。最初の引数が名前の文字列である場合にのみ使用されます。**options
key=value
いいえ
プロバイダーの構成項目 (例:
endpoint、api_key)。最初の引数が名前の文字列で、プロバイダーが渡されない場合にのみ使用されます。最初の引数がopenai-compat、dashscope、またはtritonの場合、対応するプロバイダーが作成されます。それ以外の場合は、GenericProviderが作成されます。 -
戻り値
なし。
-
例
# オプション 1:ModelProvider インスタンスを直接渡す pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) # オプション 2:カスタム名 + ModelProvider インスタンス (複数登録) pf.set_model_provider("chat", pf.OpenAICompatProvider(task="chat/completions")) pf.set_model_provider("embedding", pf.OpenAICompatProvider(task="embeddings")) # オプション 3:名前 + キーワード引数 pf.set_model_provider("dashscope", task="chat/completions")
汎用呼び出し
-
API 名:predict。
DataFrame.llm.predict( *input_cols, provider=None, model=None, output_type=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:汎用的なモデル推論を実行します。指定された入力列をモデルに送信し、モデルの出力列を DataFrame に追加します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
*input_cols
String
はい
モデル入力として使用される列の名前。複数の列を渡すことができます。
provider
String
いいえ
プロバイダー名。指定しない場合、デフォルトのプロバイダーが使用されるか、1 つだけ登録されている場合はそのプロバイダーが使用されます。
model
String
いいえ
モデル名 (例:
"qwen3.6-plus")。output_type
String / DataType
いいえ
出力型。デフォルトでは、
output STRINGが追加されます。SQL 型文字列または非構造体のDataTypeを渡すと、outputという名前の単一の列が追加されます。DataType.struct({...})を渡すと、カスタム列が追加されます。cache_table
String
いいえ
予測結果をキャッシュするために使用される Fluss テーブルの名前。
cache_key
String / List[String]
いいえ
キャッシュキー列。このパラメーターを設定する場合、
cache_tableも設定する必要があります。config
Dict
いいえ
ランタイム構成オプション。
**kwargs
key=value
いいえ
この呼び出しのモデルプロバイダーパラメーターをオーバーライドします。例えば、
OpenAICompatProviderまたはDashScopeProviderを使用する場合、content_type、system_prompt、user_prompt、temperatureなどを渡すことができます。 -
戻り値
モデル出力が追加された DataFrame。デフォルトでは
output列に追加されます。 -
例
テキスト呼び出し
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) questions = pf.from_dict({ "id": [1, 2], "question": ["What is Flink?", "What is stream processing?"], }) # デフォルトの出力列:output (STRING) df = questions.llm.predict("question", model="qwen3.6-plus") # カスタム出力列:answer (STRING) df = questions.llm.predict( "question", model="qwen3.6-plus", output_type=pf.DataType.struct({ "answer": pf.DataType.string() }))マルチモーダル呼び出し
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) tickets = pf.from_dict({ "question": ["Is this item damaged?"], "image_url": ["https://example.com/images/order-1001.jpg"], }) # マルチモーダル入力列:1 つのテキスト列と 1 つの画像列 df = tickets.llm.predict( "question", "image_url", model="qwen3.6-plus", content_type=["TEXT", "IMAGE_URL"], )結果のキャッシュ
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = questions.llm.predict( "question", model="qwen3.6-plus", cache_table="`fluss-catalog`.default_database.predict_cache", cache_key=["id", "question"])説明Flink AI Service (組み込みモデル) のみがマルチモーダル呼び出しと結果のキャッシュをサポートします。
テキスト
テキスト分類
-
API 名:ai_classify。
DataFrame.llm.ai_classify( input_col, labels, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:テキストを指定されたラベルのいずれかに分類します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
labels
List[String]
はい
分類ラベルのリスト (例:
["positive", "negative", "neutral"])。汎用呼び出しで他のパラメーターの説明をご参照ください。
説明ai_classifyなどのドメイン固有の AI 関数は、組み込みのシステムプロンプトを使用します。モデルプロバイダーで構成されたシステムプロンプトは有効になりません。 -
戻り値
分類結果と信頼度スコアを含む
category(STRING) とconfidence(DOUBLE) 列が追加された DataFrame。 -
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"review": ["Great product", "Terrible, do not buy", "It is okay"]}) df = df.llm.ai_classify("review", labels=["positive", "negative", "neutral"], model="qwen3.6-plus")
感情分析
-
API 名:ai_sentiment。
DataFrame.llm.ai_sentiment( input_col, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:入力テキストに対して感情分析を実行します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
感情分析の結果を含む、以下の列が追加された DataFrame:
-
score(DOUBLE):-1.0 から 1.0 の範囲の感情分析スコア。 -
label(STRING):「positive」、「negative」、または「neutral」のいずれか。 -
confidence(DOUBLE):信頼度スコア。
-
-
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"comment": ["This feature is amazing!", "Broke after one month"]}) df = df.llm.ai_sentiment("comment", model="qwen3.6-plus")
情報抽出
-
API 名:ai_extract。
DataFrame.llm.ai_extract( input_col, schema, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:指定された JSON スキーマに基づいて、テキストから構造化情報を抽出します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
schema
String
はい
抽出するフィールドを記述する JSON スキーマ文字列 (例:
'{"name":"STRING", "phone":"STRING"}')。汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
抽出された構造化情報を含む
extracted_json(STRING) 列が追加された DataFrame。 -
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"text": ["John Smith, male, 28 years old, phone ***-****"]}) schema = '{"name": "STRING", "age": "INTEGER", "phone": "STRING"}' df = df.llm.ai_extract("text", schema=schema, model="qwen3.6-plus")
テキスト翻訳
-
API 名:ai_translate。
DataFrame.llm.ai_translate( input_col, source_lang, target_lang, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:テキストをソース言語からターゲット言語に翻訳します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
source_lang
String
はい
ソース言語コード (例:
"zh"、"en"、"auto"(自動検出))。サポートされている言語:auto、zh、en、ja、ko、fr、de、es、ru、ar、pt。target_lang
String
はい
ターゲット言語コード。この値は
"auto"にはできません。汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
翻訳されたテキストと検出された言語を含む
translated_text(STRING) とdetected_language(STRING) 列が追加された DataFrame。 -
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"text": ["Hello World", "How are you?"]}) df = df.llm.ai_translate("text", source_lang="en", target_lang="zh", model="qwen3.6-plus")
テキスト要約
-
API 名:ai_summarize。
DataFrame.llm.ai_summarize( input_col, max_length, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:テキストを指定された最大長に要約します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
max_length
int
はい
要約の最大文字数。0 より大きい必要があります。
汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
要約テキストを含む
summary(STRING) 列が追加された DataFrame。 -
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"article": ["This is a long article with lots of content..."]}) df = df.llm.ai_summarize("article", max_length=100, model="qwen3.6-plus")
データマスキング
-
API 名:ai_mask。
DataFrame.llm.ai_mask( input_col, entities, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:テキスト内の機密情報をマスキングします。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
entities
List[String]
はい
マスキングするエンティティの種類のリスト (例:
["name", "phone"])。汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
マスキングされたテキストと検出されたエンティティを含む
masked_text(STRING) とdetected_entities(ARRAY<STRING>) 列が追加された DataFrame。 -
例
pf.set_model_provider(pf.OpenAICompatProvider(task="chat/completions")) df = pf.from_dict({"text": ["Please contact John Smith at 555-0123"]}) df = df.llm.ai_mask("text", entities=["name", "phone"], model="qwen3.6-plus")
ベクター
ベクトル検索
-
vector_searchメソッドは VVR 11.8 以降でサポートされています。 -
このトピックで説明されている
output_columnsの結果列の射影とマッピングベースのエイリアス機能は、VVR 11.9 以降でサポートされています。 -
VVR 11.8 では、
output_columnsはすべての出力列の名前を変更するだけです。非集約モードでは、ベクトルテーブルの列数とscoreに相当する数の名前を提供する必要があります。集約モードでは、外部配列列に 1 つの名前を提供する必要があります。VVR 11.9 以降にアップグレードした後は、このパラメーターをこのトピックで説明されている射影セマンティクスに合わせて調整してください。
-
API 名:vector_search。
DataFrame.llm.vector_search( search_source, column_to_search, column_to_query, top_k, *, agg=False, output_columns=None, config=None, ignore_empty=False ) -
説明:現在の DataFrame のクエリベクターを使用して、ベクトル検索ソースで最も類似したレコードを検索し、検索結果を現在の DataFrame に追加します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
search_source
DataFrame
はい
ベクトル検索ソースの DataFrame (例:
read_milvusを使用して読み取られたベクトルコレクション)。column_to_search
String
はい
検索ソース内のベクトル列の名前。
column_to_query
String / Expression
はい
現在の DataFrame 内のクエリベクトル列名または列式。
top_k
int
はい
入力レコードごとに返される最も類似したレコードの数。
agg
Boolean
いいえ
上位 K 件の結果を単一の配列列に集約するかどうかを指定します。デフォルト値:
False。output_columns
List[String] / Mapping[String, String]
いいえ
返すベクトル検索結果の列。リストは列を選択し、元の名前を保持します。
agg=Falseの場合、マッピングも使用でき、キーが列を選択し、値が出力エイリアスになります。search_sourceのトップレベルの列と、エンジンによって生成されたscore列を選択できます。デフォルト値:None。これはすべての結果列を選択します。config
Dict
いいえ
ランタイムのベクトル検索構成。
ignore_empty
Boolean
いいえ
検索結果がない入力行をドロップするかどうかを指定します。デフォルト値:
False。 -
戻り値
元の列と追加されたベクトル検索結果の列を含む DataFrame オブジェクト。
output_columnsが合計 M 個の結果列を選択すると仮定します。-
`output_columns` を指定しない場合、
search_sourceのすべての列と、DOUBLE型の検索スコア列が返されます。デフォルトでは、スコア列の名前はscoreです。スコアの意味とソート方向は、コネクタとその距離メトリックに依存します。Milvus の場合、値をsearch_metricと一緒に解釈します。 -
`agg=False` の場合、各入力行は最大で `top_k` 個の出力行を生成し、M 個の列が追加されます。リストは宣言された順序で列を選択し、元の名前を保持します。マッピングはキーの宣言順序で列を選択し、対応する値を出力エイリアスとして使用します。
-
`agg=True` の場合、検索結果を持つ各入力は最大で 1 つの出力行を生成し、
ARRAY<ROW<...>>型のsearch_results列が追加されます。デフォルトではignore_empty=Falseはすべての入力レコードを保持し、入力に検索結果がない場合search_resultsはNULLになります。ignore_empty=Trueを設定すると、検索結果のない入力は結果に表示されません。各配列要素は、M 個の選択されたフィールドを含む ROW です。このモードは、内部フィールドを射影するためのリストのみをサポートし、マッピングはサポートしません。呼び出し後、`rename_columns("search_results", "new_name")` を呼び出して外部配列列の名前を変更します。 -
`output_columns` は空にできず、列名もエイリアスも空または重複にすることはできません。選択された列は、`search_source` のトップレベルの列または生成された類似度列でなければなりません。返される順序は、リストまたはマッピングキーの宣言順序と一致します。
-
`search_source` に `score` や `score0` などの列が既にある場合、エンジンは生成された類似度列に `score2` などの一意の名前を選択します。この場合、`output_columns` で実際の列名を使用します。
-
非集約モードでは、マッピングを使用して、競合しない結果列のエイリアスを指定できます。リストまたはデフォルトの出力が現在の DataFrame の既存の列と競合する場合、`vector_search` を呼び出す前に `rename_columns` で競合する列の名前を変更します。集約モードでは、現在の DataFrame に既に `search_results` 列がある場合、まずその列の名前を変更します。
-
ベクトル検索ソースのコネクタが列プルーニングのプッシュダウンをサポートしている場合、エンジンは `output_columns` のプルーニングをコネクタにプッシュダウンします。プッシュダウンがサポートされているかどうかに関わらず、最終結果には選択されたベクトル検索結果の列のみが含まれます。
-
-
VVR 11.8 からのアップグレードに関する注意
VVR 11.8 では、
output_columnsはすべての出力列のエイリアスを意味します。VVR 11.9 以降では、選択する結果列を意味します。例えば、VVR 11.8 用に書かれた以下の呼び出し:output_columns=["match_doc_id", "match_embedding", "match_title", "match_score"]アップグレード後、代わりにマッピングを使用します。
output_columns={ "doc_id": "match_doc_id", "embedding": "match_embedding", "title": "match_title", "score": "match_score", } -
例
from array import array import pyflink.dataframe as pf from pyflink.dataframe import DataType query_df = pf.from_dict({ "query_id": [1], "query_embedding": [array("f", [0.1, 0.2, 0.3])], }) documents = pf.read_milvus( endpoint="http://milvus.example.com", username="${secret_values.milvus_user}", password="${secret_values.milvus_password}", database_name="commerce", collection_name="support_docs", schema={ "doc_id": DataType.int64(), "embedding": DataType.list(DataType.float32()), "title": DataType.string(), }, search_metric="COSINE", ) selected_docs = query_df.llm.vector_search( documents, column_to_search="embedding", column_to_query="query_embedding", top_k=3, output_columns=["doc_id", "title", "score"], )リストは結果列を射影し、元の名前を保持します。同時にエイリアスを指定するには、マッピングを使用します。
matched_docs = query_df.llm.vector_search( documents, column_to_search="embedding", column_to_query="query_embedding", top_k=3, output_columns={ "doc_id": "match_doc_id", "title": "match_title", "score": "match_score", }, )集約モードでは、リストは各配列要素のフィールドを選択し、外部配列列は常に
search_resultsという名前になります。外部配列列の名前を変更するには、rename_columnsを呼び出します。matched_docs_agg = query_df.llm.vector_search( documents, column_to_search="embedding", column_to_query="query_embedding", top_k=3, agg=True, output_columns=["doc_id", "title", "score"], ).rename_columns("search_results", "matches")
テキスト埋め込み
-
API 名:ai_embed。
DataFrame.llm.ai_embed( input_col, dimension=1024, *, provider=None, model=None, cache_table=None, cache_key=None, config=None, **kwargs ) -
説明:テキストのベクトル埋め込みを生成します。
-
パラメーター
パラメーター名
パラメーターの型
必須
説明
input_col
String / Expression
はい
入力テキストの列名または列式。
dimension
int
いいえ
ベクトルのディメンション。デフォルト値:1024。
汎用呼び出しで他のパラメーターの説明をご参照ください。
-
戻り値
生成されたベクトルを含む
embedding(ARRAY<FLOAT>) 列が追加された DataFrame。 -
例
pf.set_model_provider("embedding", pf.OpenAICompatProvider(task="embeddings")) df = pf.from_dict({"text": ["stream processing with Flink", "real-time analytics"]}) df = df.llm.ai_embed("text", dimension=1024, provider="embedding", model="text-embedding-v4")
マルチモーダル演算子
DataFrame API を使用すると、マルチモーダル演算子を呼び出して、画像、音声、ビデオなどのマルチモーダルデータを読み取り、処理、分析できます。詳細については、「マルチモーダル演算子」をご参照ください。API ドキュメントについては、「マルチモーダル式」をご参照ください。
環境と構成
|
API |
説明 |
|
グローバルな TableEnvironment の設定 |
|
|
現在の TableEnvironment の取得 |
|
|
TableEnvironment の取得または自動作成 |
|
|
DataFrame 構成オプションの設定 |
|
|
DataFrame 構成オプションの読み取り |
実行時パラメーターの変更
-
API 名:DataFrameConfig.set。
DataFrameConfig.set( key, value ) -
説明: DataFrame API ジョブの実行時パラメーターを変更します。 次の例に示すように、Python ジョブパラメーターおよびTable API パラメーターを構成できます。 構成されたパラメーターは、デフォルトの
TableEnvironmentに適用され、ジョブの [実行時パラメーター構成] にある同じ名前のパラメーターをオーバーライドします。 -
パラメーター
パラメーター名
パラメーターの型
必須
説明
key
String
はい
パラメーター名。
value
String
はい
パラメーター値。
-
戻り値
DataFrameConfig オブジェクト自体。メソッドチェーンをサポートします。
-
例
import pyflink.dataframe as pf pf.config.set("python.fn-execution.arrow.batch.size", "256") pf.config.set("table.exec.async-scalar.max-concurrent-operations", "10") pf.config.set("table.exec.async-lookup.timeout", "1 min")説明パラメーターが期待どおりに有効になるように、ジョブログを定義する前にパラメーターを構成してください。