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

Realtime Compute for Apache Flink:Python DataFrame API リファレンス

最終更新日:Sep 21, 2026

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

属性

schema、columns

射影/列操作

select、with_column、with_columns、drop_columns、rename_columns

フィルター

filter

集約

group_by、agg

結合/和集合

join、join_asof、join_lateral

行マッピング

map、map_batches、flat_map

展開

explode

重複排除

drop_duplicates

集合演算

union、union_all、minus、minus_all、intersect、intersect_all

制限/ページネーション

limit、offset、head

パーティション

rebalance

パイプ

pipe

Null 処理

drop_null、drop_nan、fill_null、fill_nan

出力/収集

to_table、to_pandas、collect、iter_rows、iter_batches

デバッグ/実行計画

explain

SQL

sql

式ヘルパー関数

関数

説明

col

列参照式を作成します。

lit

リテラル式を作成します。

データ型

DataType は DataFrame の列の型を記述するために使用されます。

型カテゴリ

メソッド

ブール値

DataType.bool

整数

DataType.int8、DataType.int16、DataType.int32、DataType.int64

浮動小数点/固定小数点

DataType.float32、DataType.float64、DataType.decimal

文字列

DataType.string、DataType.fixed_size_string

バイナリ

DataType.binary、DataType.fixed_size_binary

日付/時刻

DataType.date、DataType.time、DataType.timestamp、DataType.timestamp_ltz

複合

DataType.list、DataType.map、DataType.struct

特殊

DataType.null、DataType.variant

マルチモーダル

DataType.tensor、DataType.image

Null 許容性修飾子

DataType.not_null、DataType.nullable

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 ファイル

read_parquet

JSON ファイル

read_json

Kafka

read_kafka

MaxCompute (ODPS)

read_odps

Paimon

read_paimon

SLS (Simple Log Service)

read_sls

Hologres

read_hologres

Milvus

read_milvus

ビデオフレーム

read_video_frames

汎用/カスタムコネクタ

read_generic

データの書き込み

シンク

API

Parquet ファイル

write_parquet

JSON ファイル

write_json

Kafka

write_kafka

MaxCompute (ODPS)

write_odps

Paimon

write_paimon

SLS (Simple Log Service)

write_sls

Hologres

write_hologres

Milvus

write_milvus

プリント

write_print

VVR 11.9.preview.1 からサポートされています。

汎用/カスタムコネクタ

write_generic

単一のジョブで複数回のデータ書き込みが必要な場合は、create_statement_set を使用して StatementSet を作成し、すべての書き込み関数をまとめて実行します。

カタログの使用

DataFrame API はカタログを使用し、カタログテーブルからの読み取りや書き込みができます。

API

説明

create_catalog

カタログの作成

get_catalog

登録済みのカタログオブジェクトの取得

use_catalog

現在のカタログの切り替え

get_current_catalog

現在のカタログの取得

list_catalogs

カタログの一覧表示

use_database

現在のデータベースの切り替え

get_current_database

現在のデータベースの取得

list_databases

現在のカタログ内のデータベースの一覧表示

read_catalog_table

カタログテーブルを DataFrame として読み取る

write_catalog_table

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

説明

set_model_provider

モデルプロバイダーを登録します。

set_default_model_provider

デフォルトのプロバイダーを設定します (複数プロバイダーのシナリオ向け)。

list_model_providers

登録済みのプロバイダーを一覧表示します。

サポートされているプロバイダーの種類:

プロバイダー

ユースケース

OpenAICompatProvider

OpenAI、DeepSeek、DashScope (Model Studio)、およびすべての OpenAI 互換インターフェイス。

DashScopeProvider

マルチモーダル埋め込みをサポートする Alibaba Cloud DashScope。

TritonProvider

NVIDIA Triton Inference Server

GenericProvider

登録済みのカスタムモデルプロバイダーを使用します。

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

説明

set_table_environment

グローバルな TableEnvironment の設定

get_table_environment

現在の TableEnvironment の取得

get_or_create_table_environment

TableEnvironment の取得または自動作成

config.set(key, value)

DataFrame 構成オプションの設定

config.get(key)

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")
    説明

    パラメーターが期待どおりに有効になるように、ジョブログを定義する前にパラメーターを構成してください。