全部產品
Search
文件中心

Realtime Compute for Apache Flink:功能概覽

更新時間:Jul 30, 2026

本文以一個電商訂單分析和客服問題處理的完整情境為例,為您介紹 DataFrame API 的功能特性。

構建 DataFrame 對象

使用本機資料

在進行開發調試時,可以用 Python 字典、記錄列表或 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

# User table
users = pf.from_dict({
    "user_id": [1, 2, 3],
    "name": ["Alice", "Bob", "Charlie"],
    "city": ["Hangzhou", "Shanghai", "Beijing"],
})

# Order fact table
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"],
)

# User behavior event table
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"],
)

# User profile table
profiles = pf.from_records(
    [
        (1, "gold"),
        (2, "silver"),
        (3, "gold"),
    ],
    schema=["user_id", "member_level"],
)

# You can also build from 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 參數為您使用連接器的 identifier,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",
    },
)

讀取 Catalog 表

如果資料已經註冊在 Flink Catalog 中,可以使用 read_catalog_table 讀取 Catalog 表。

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]

缺失值和 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"]
)

展開

如果欄位中包含數組、Map 或 Multiset,可以使用 explode 將一行展開為多行:

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 可用於執行維表 join,右表資料需要從維表資料來源讀取。

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)將一行輸入資料處理為一個結果值。DataFrame API 支援同步/非同步、逐行/向量化等多種 UDF 形態。UDF 可以在 with_column 或 with_columns 中使用,將結果添加為新列;也可以在 map 或 map_batches 中使用,對整行資料做變換。

  • 自訂表格值函數(udtf)將一行輸入展開為任意行輸出,每行輸出可包括一列或多列結果。可以通過 join_lateral 或 flat_map 調用。

下面的樣本使用使用者自訂函數處理資料。完整使用方法請參見使用者自訂函數。

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 註冊模型 Provider,再在 DataFrame 上調用通用預測或內建 AI 函數。

配置 Provider

以下樣本使用 Flink AI 服務提供的內建模型,僅需指定 task 參數,無需配置 endpoint 和 api_key。需開通 Flink AI 服務後方可使用內建模型,詳情請參見 Flink AI服務(內建模型)。

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 函數。

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

# sentiment analysis
sentiment = questions.llm.ai_sentiment("question", provider="chat", model="qwen3.6-plus")

# information extraction, schema described as JSON string
extracted = questions.llm.ai_extract(
    "question",
    '{"order_id":"STRING","intent":"STRING"}',
    provider="chat",
    model="qwen3.6-plus",
)

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

# mask sensitive data
masked = questions.llm.ai_mask(
    "question",
    entities=["PERSON", "PHONE", "EMAIL"],
    provider="chat",
    model="qwen3.6-plus",
)

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

上述 AI 函數還支援通過 cache_table 和 cache_key 參數,使用 Fluss 表緩衝重複推理結果,降低 token 開銷。詳情參見通用調用。

向量搜尋

您可以通過 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 文檔請參見 Multimodal Expressions。

下面的例子從圖片 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 參數為您使用連接器的 identifier,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",
    },
)

寫入 Catalog 表

如果目標表已經存在於 Catalog 中,也可以直接使用 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()