このトピックでは、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()