本文以一個電商訂單分析和客服問題處理的完整情境為例,為您介紹 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()