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

Realtime Compute for Apache Flink:機能概要

最終更新日:Aug 04, 2026

このトピックでは、e コマースの注文分析とカスタマーサポートのチケット処理というエンドツーエンドのシナリオを通して、DataFrame API の機能を紹介します。

DataFrameオブジェクトの構築

ローカルデータの使用

開発およびデバッグ中に、Python の dict、レコードのリスト、または Pandas DataFrame から DataFrame を迅速に構築できます。このトピック全体で使用される DataFrame は、ユーザーテーブル users、注文テーブル orders、イベントテーブル events、およびユーザープロファイルテーブル profiles として、以下のように構築します。

import pandas as pd
import pyflink.dataframe as pf
from pyflink.dataframe import col, lit, udf, udtf, DataType

# ユーザーテーブル
users = pf.from_dict({
    "user_id": [1, 2, 3],
    "name": ["Alice", "Bob", "Charlie"],
    "city": ["Hangzhou", "Shanghai", "Beijing"],
})

# 注文ファクトテーブル
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"],
)

# ユーザー行動イベントテーブル
events = pf.from_records(
    [
        {"event_id": "e1", "user_id": 1, "event": "login"},
        {"event_id": "e2", "user_id": 2, "event": "pay"},
    ],
    schema=["event_id", "user_id", "event"],
)

# ユーザープロファイルテーブル
profiles = pf.from_records(
    [
        (1, "gold"),
        (2, "silver"),
        (3, "gold"),
    ],
    schema=["user_id", "member_level"],
)

# pandas DataFrame からも構築できます
pdf = pd.DataFrame({"id": [1, 2], "score": [0.8, 0.95]})
df_from_pandas = pf.from_pandas(pdf)

外部データの読み取り

組み込みの読み取り関数の使用

本番環境では、DataFrame API の組み込みの読み取り関数を使用して、一般的なデータソースから読み取ることができます。

kafka_events = pf.read_kafka(
    "broker-1:9092,broker-2:9092",
    topic="user_events",
    group_id="df-api-demo",
    startup_mode="earliest-offset",
    format="json",
    schema={
        "user_id": DataType.int64(),
        "event": DataType.string(),
        "event_time": DataType.timestamp(3),
        "payload": DataType.string(),
    },
)

その他の読み取り関数については、「データインジェスト」をご参照ください。

説明

外部システムから読み取る場合、通常、フィールドタイプを明示的に宣言する必要があります。DataType は、整数、浮動小数点数、文字列、ブール値、日時、複合型などの一般的な型を提供します。

カスタムコネクタの使用

read_generic 関数を使用して、他のサポートされているコネクタやカスタムコネクタからデータを読み取ることができます。

connector パラメーターは使用するコネクタの識別子を示し、options はコネクタの WITH パラメーターを渡します。

source = pf.read_generic(
    "datagen",
    schema={
        "id": DataType.int64(),
        "name": DataType.string(),
    },
    options={
        "number-of-rows": "1000",
        "fields.id.kind": "sequence",
        "fields.id.start": "1",
        "fields.id.end": "1000",
    },
)

カタログテーブルの読み取り

データが既に Flink カタログに登録されている場合は、read_catalog_table を使用してカタログテーブルから読み取ることができます。

pf.create_catalog(
    "lake",
    {
        "type": "paimon",
        "warehouse": "oss://my-bucket/warehouse",
    },
)

pf.use_catalog("lake")
pf.use_database("commerce")

historical_orders = pf.read_catalog_table("orders")

データ処理

ほとんどの DataFrame 変換は新しい DataFrame を返すため、連鎖に適しています。一般的な操作には、射影、フィルタリング、派生列、列の削除、名前の変更、集約、結合、重複排除、セット操作、配列のエクスプロード、欠損値の処理などがあります。

射影、フィルタリング、派生列

次の paid_orders は、注文テーブル orders から支払い済みの注文をフィルタリングし、金額によってバケット化します:

paid_orders = (
    orders
    .select("order_id", "user_id", "amount", "status")
    .filter("status = 'PAID' AND amount > 0")
    .with_columns(
        amount_yuan=col("amount"),
        amount_level=(
            (col("amount") >= 100).then("high", "normal")
        ),
    )
    .drop("status", "amount")
    .rename({"amount_yuan": "amount"})
)

添字構文を使用して列を参照または選択することもできます:

amount_expr = paid_orders["amount"]
simple_view = paid_orders[["order_id", "user_id", "amount"]]
large_orders = paid_orders[col("amount") > 100]

null値とNaN値の処理

drop_null / fill_null は SQL の NULL を処理し、drop_nan / fill_nan は浮動小数点列の NaN を処理します。

次の clean_orders は、品質クレンジング後の注文テーブルを表します:

clean_orders = (
    orders
    .drop_null(subset=["order_id", "user_id"])
    .fill_null("UNKNOWN", subset=["status"])
    .fill_nan(0.0, subset=["amount"])
)

重複排除

drop_duplicates を使用して、すべてのフィールド、または選択したフィールドのサブセットに基づいて重複を排除します。

unique_events = events.drop_duplicates(
    subset=["event_id"]
)

エクスプロード

フィールドに配列、マップ、またはマルチセットが含まれている場合、explode を使用して 1 行を複数行に展開できます:

order_items = pf.from_records(
    [
        (1001, ["phone", "case"]),
        (1002, ["book"]),
    ],
    schema=["order_id", "items"],
)

item_rows = order_items.explode("items", "item")

セット操作

構造的に互換性のある複数の DataFrame は、union、union_all、intersect、intersect_all、minus、および minus_all を使用して結合できます。

説明

ストリーミングモードでは、union_all のみがサポートされます。

all_paid_orders = (
    paid_orders.select("order_id", "user_id", "amount")
    .union_all(historical_orders.select("order_id", "user_id", "amount"))
    .drop_duplicates(subset=["order_id"])
)

集約

group_by を agg と一緒に使用して、グループごとの集約を実行します。グループ化列は、結果の DataFrame に保持されます。

次の user_summary は、ユーザーごとに集約された注文メトリクステーブルで、注文数、合計金額、平均金額が含まれます:

user_summary = (
    clean_orders
    .filter("status = 'PAID'")
    .group_by("user_id")
    .agg(
        order_count=col("order_id").count,
        total_amount=col("amount").sum,
        avg_amount=col("amount").avg,
    )
)

テーブル全体で集約することもできます:

overall = clean_orders.agg(
    order_count=col("order_id").count,
    total_amount=col("amount").sum,
)

ジョイン

join は、inner、left、right、full、および outer のジョインタイプをサポートします。

次の例では、ユーザーテーブル users をイベントテーブル events にジョインして、ユーザーの名前、都市、イベントを含む enriched ワイドテーブルを生成します:

enriched = (
    users
    .join(
        events,
        on="user_id",
        how="left",
    )
    .select(
        "user_id",
        "name",
        "city",
        "event_id",
        "event",
    )
)

join_asof を使用して、ディメンションテーブルジョインを実行できます。右側のデータは、ディメンションテーブルソースからのものである必要があります。

dim_users = pf.read_hologres(
    endpoint="hgpostcn-cn-xxx-cn-hangzhou.hologres.aliyuncs.com:80",
    db_name="commerce",
    table_name="dim_user_profile",
    username="${secret_values.hg_user}",
    password="${secret_values.hg_password}",
    schema={
        "user_id": DataType.int64(),
        "member_level": DataType.string(),
        "city": DataType.string(),
    },
    primary_key="user_id",
    binlog=False,
    cache="LRU",
)

events_with_profile = events.join_asof(
    dim_users,
    by="user_id",
    how="left",
)

SQLクエリの使用

sql を使用して SELECT クエリを実行できます。SQL は、手動で登録することなく、現在のスコープ内の DataFrame 変数とユーザー定義関数を参照できます。

次の例では、SQL を使用して、ユーザーごとの注文サマリー user_summary から高額支出ユーザーを選択し、ユーザープロファイルテーブル profiles とジョインします:

top_users = pf.sql("""
    SELECT
        user_summary.user_id,
        profiles.member_level,
        user_summary.order_count,
        user_summary.total_amount
    FROM user_summary
    LEFT JOIN profiles
        ON user_summary.user_id = profiles.user_id
    WHERE user_summary.total_amount >= 200
""")

ユーザー定義関数

標準の式、集約、または SQL でビジネスロジックを表現できない場合は、ユーザー定義関数を使用できます:

  • ユーザー定義スカラー関数 (udf) は、1 つの入力行を処理して単一の結果値を生成します。DataFrame API は、同期/非同期および行単位/ベクトル化 UDF をサポートしています。UDF は with_column または with_columns で使用して結果を新しい列として追加したり、map または map_batches で行全体を変換したりできます。

  • ユーザー定義テーブル関数 (udtf) は、1 つの入力行を任意の数の出力行に展開し、各行は 1 つ以上の列を持ちます。UDTF は join_lateral または flat_map で呼び出すことができます。

次の例では、UDF を使用してデータを処理します。完全な使用方法については、「ユーザー定義関数」をご参照ください。

from typing import Iterator, TypedDict

@udf
def normalize_status(status: str) -> str:
    if status is None:
        return "UNKNOWN"
    return status.strip().upper()

orders_with_status = orders.with_column(
    "status_norm", normalize_status(col("status"))
)

@udtf
def review_reasons(status: str, amount: float) -> Iterator[str]:
    if status is None:
        yield "missing_status"
    elif status == "CREATED":
        yield "pending_payment"
    if amount is not None and amount >= 100:
        yield "high_value_order"

order_review_reasons = orders.join_lateral(
    review_reasons(col("status"), col("amount")).alias("review_reason"),
)

AI / LLM

DataFrame API は、df.llm を通じて AI/LLM 機能を提供します。まず set_model_provider でモデルプロバイダーを登録し、次に DataFrame で汎用呼び出しまたは組み込み AI 関数を呼び出します。

プロバイダーの設定

次の例では、Flink AI Service が提供する組み込みモデルを使用します。task パラメーターのみが必要で、endpoint や api_key を設定する必要はありません。組み込みモデルを使用する前に、Flink AI Service を有効にする必要があります。詳細については、「Flink AI Service (組み込みモデル)」をご参照ください。

pf.set_model_provider("chat", pf.OpenAICompatProvider(task="chat/completions"))
pf.set_model_provider("embedding", pf.OpenAICompatProvider(task="embeddings"))

汎用呼び出し

predict は汎用的な LLM の呼び出しに使用され、出力フィールドは output_type で指定できます。

たとえば、次の questions はカスタマーサポートのチケットの質問を表し、answers はモデルの回答を含む output 列を追加した結果テーブルです。

questions = pf.from_records(
    [
        (1, "Why hasn't my order been shipped yet? My order number is 1. Charlie"),
        (2, "How do I apply for a refund? My email is bob@example.com. Order number: 2."),
    ],
    schema=["ticket_id", "question"],
)

answers = questions.llm.predict(
    "question",
    provider="chat",
    model="qwen3.6-plus",
    system_prompt="You are a helpful assistant.",
    user_prompt="Answer the customer support question concisely. Return only the answer.",
    temperature=0.2,
    config={"max-concurrent-operations": "20"},
)

組み込みAI関数

一般的なタスクには、組み込み AI 関数を使用できます。組み込み AI 関数は、元の DataFrame に結果列を追加します。詳細については、「AI / LLM 関数」をご参照ください。

# 分類
classified = questions.llm.ai_classify(
    "question",
    labels=["logistics", "refund", "account", "other"],
    provider="chat",
    model="qwen3.6-plus",
)

# 感情分析
sentiment = questions.llm.ai_sentiment("question", provider="chat", model="qwen3.6-plus")

# 情報抽出、スキーマはJSON文字列として記述します
extracted = questions.llm.ai_extract(
    "question",
    '{"order_id":"STRING","intent":"STRING"}',
    provider="chat",
    model="qwen3.6-plus",
)

# 要約
summaries = questions.llm.ai_summarize(
    "question",
    max_length=50,
    provider="chat",
    model="qwen3.6-plus",
)

# 機密データのマスキング
masked = questions.llm.ai_mask(
    "question",
    entities=["PERSON", "PHONE", "EMAIL"],
    provider="chat",
    model="qwen3.6-plus",
)

# 埋め込み
embeddings = questions.llm.ai_embed(
    "question",
    dimension=1024,
    provider="embedding",
    model="text-embedding-v4",
)

上記の AI 機能は、cache_table と cache_key パラメーターを介して Fluss テーブルに繰り返しの推論結果をキャッシュすることもサポートしており、トークンコストを削減します。詳細については、「汎用呼び出し」をご参照ください。

ベクトル検索

vector_search を使用して、Milvus などのベクトルデータベースに接続し、ベクトル検索を実行できます。

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",
)

matched_docs = embeddings.llm.vector_search(
    documents,
    column_to_search="embedding",
    column_to_query="embedding",
    top_k=3,
    output_columns=["doc_id", "doc_embedding", "doc_title", "score"],
)

マルチモーダルオペレーター

DataFrame API を使用すると、マルチモーダルオペレーターを便利に呼び出して、画像、音声、動画などのマルチモーダルデータを読み取り、分析し、変換できます。詳細については、「マルチモーダルオペレーター」をご参照ください。API リファレンスについては、「マルチモーダル式」をご参照ください。

次の例では、URL から画像コンテンツを読み取り、有効性チェック、デコード、メタデータ抽出、品質スコアリング、サムネイル生成を実行します。

images = pf.from_dict({
    "image_id": ["1", "2"],
    "image_url": [
        "oss://bucket/products/1.jpg",
        "oss://bucket/products/2.jpg",
    ],
})

image_analysis = (
    images
    .with_column("image_bytes", col("image_url").fetch_content())
    .with_columns(
        image_valid=col("image_bytes").image.is_valid(
            pixel_limit=40_000_000
        ),
        metadata=col("image_bytes").image.metadata(),
    )
    .filter("image_valid")
    .with_column(
        "image",
        col("image_bytes").image.decode(
            on_error="null", mode="RGB", pixel_limit=40_000_000
        ),
    )
    .drop_null(subset=["image"])
    .with_columns(
        quality_score=col("image").image.quality_score(),
        thumbnail=(
            col("image")
            .image.resize(width=512, height=512)
            .image.encode(format="JPEG", quality=85)
        ),
    )
)

データシンク

DataFrame の書き込みメソッドは、現在の DataFrame をターゲットシステムに書き込むために使用します。

組み込みシンク関数の使用

DataFrame API の組み込みシンク関数を使用して、一般的なシステムにデータを書き込むことができます。

enriched.write_kafka(
    "broker-1:9092,broker-2:9092",
    topic="user_summary",
    format="json",
    key_format="json",
    key_fields=["user_id"],
    delivery_guarantee="at-least-once",
)

その他のシンク関数については、「データシンク」をご参照ください。

カスタムコネクタの使用

write_generic 関数を使用して、他のサポートされているコネクタやカスタムコネクタを介してデータを書き込むことができます。

connector パラメーターは使用するコネクタの識別子を示し、options はコネクタの WITH パラメーターを渡します。

enriched.write_generic(
    "blackhole",
    options={
        "sink.parallelism": "4",
    },
)

コネクタが主キー制約を必要とする場合は、primary_key を渡します:

enriched.write_generic(
    "my-custom-sink",
    primary_key="user_id",
    options={
        "endpoint": "example.com:9000",
        "format": "json",
    },
)

カタログテーブルへの書き込み

ターゲットテーブルが既にカタログに存在する場合は、write_catalog_table を使用して直接書き込むこともできます:

matched_docs.write_catalog_table(
    "lake.commerce.support_ticket_matches",
    overwrite=True,
)

複数のターゲットへの書き込み

同じ計算パイプラインを複数のターゲットに書き込む必要がある場合は、create_statement_set を使用して StatementSet を作成し、それを各書き込みメソッドの statement_set パラメーターに渡して、まとめて実行します:

statement_set = pf.create_statement_set()

enriched.write_kafka(
    "broker-1:9092,broker-2:9092",
    topic="user_summary",
    format="json",
    statement_set=statement_set,
)
enriched.write_json(
    "oss://my-bucket/user_summary/",
    statement_set=statement_set,
)

statement_set.execute()