Feature overview

Updated at:

This topic walks through the features of the DataFrame API using an end-to-end scenario of e-commerce order analysis and customer-support ticket handling.

Build DataFrame objects

Use local data

During development and debugging, you can quickly build a DataFrame from a Python dict, a list of records, or a Pandas DataFrame. The DataFrames used throughout this topic are built as follows: user table users, order table orders, event table events, and user profile table 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)

Read external data

Use built-in read functions

In production, you can use the built-in read functions of the DataFrame API to read from common data sources.

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(),
    },
)

For more read functions, see Data ingestion.

Note

When reading from external systems, field types usually need to be declared explicitly. DataType provides common types, including integers, floats, strings, booleans, date-time, and composite types.

Use a custom connector

You can use the read_generic function to read data from other supported connectors or custom connectors.

The connector parameter is the identifier of the connector to use, and options passes the WITH parameters of the connector.

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

Read a Catalog table

If your data is already registered in a Flink Catalog, you can use read_catalog_table to read from Catalog tables.

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

Data processing

Most DataFrame transformations return a new DataFrame, which makes them well suited to chaining. Common operations include projection, filtering, derived columns, dropping columns, renaming, aggregation, joining, deduplication, set operations, exploding arrays, and handling missing values.

Projection, filtering, and derived columns

The following paid_orders filters paid orders from the order table orders and buckets them by amount:

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

You can also use subscript syntax to reference or select columns:

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

Handle null and NaN values

drop_null / fill_null handle SQL NULL, and drop_nan / fill_nan handle NaN in floating-point columns.

The following clean_orders represents an order table after quality cleansing:

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

Deduplication

Use drop_duplicates to deduplicate by all fields or a chosen subset.

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

Explode

When a field contains an array, map, or multiset, you can use explode to expand one row into multiple rows:

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

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

Set operations

Multiple structurally compatible DataFrames can be combined with union, union_all, intersect, intersect_all, minus, and minus_all.

Note

In streaming mode, only union_all is supported.

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

Aggregation

Use group_by together with agg to perform grouped aggregation. Grouping columns are retained in the resulting DataFrame.

The following user_summary is an order-metrics table aggregated by user, containing order count, total amount, and average amount:

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

You can also aggregate across the entire table:

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

Join

join supports the inner, left, right, full, and outer join types.

The following example joins the user table users with the event table events to produce an enriched wide table containing the user's name, city, and event:

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

join_asof can be used to perform a dimension-table join. The right-side data must come from a dimension-table source.

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

Use SQL queries

You can run SELECT queries with sql. The SQL can reference DataFrame variables and user-defined functions in the current scope without manual registration.

The following example uses SQL to select high-spend users from the per-user order summary user_summary and join with the user profile table 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
""")

User-defined functions

When standard expressions, aggregation, or SQL cannot express the business logic, you can use user-defined functions:

  • A user-defined scalar function (udf) processes one input row into a single result value. The DataFrame API supports sync/async and row-by-row/vectorized UDFs. UDFs can be used in with_column or with_columns to add the result as a new column, or in map or map_batches to transform whole rows.

  • A user-defined table function (udtf) expands one input row into any number of output rows, each with one or more columns. UDTFs can be invoked with join_lateral or flat_map.

The following example uses UDFs to process data. For the complete usage, see User-defined functions.

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

The DataFrame API exposes AI/LLM capabilities through df.llm. First register a model provider with set_model_provider, and then call generic invocation or a built-in AI function on the DataFrame.

Configure a provider

The following example uses the built-in models provided by Flink AI Service. Only the task parameter is required; you do not need to configure endpoint or api_key. Flink AI Service must be enabled before you can use the built-in models. For more information, see Flink AI Service (built-in models).

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

Generic invocation

predict is used for generic LLM invocation. The output field can be specified with output_type.

For example, the following questions represents customer-support ticket questions, and answers is the result table with an appended output column containing the model's answers.

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

Built-in AI functions

You can use built-in AI functions for common tasks. Built-in AI functions append result columns to the original DataFrame; for details, see AI / LLM functions.

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

The AI functions above also support caching repeated inference results in a Fluss table through the cache_table and cache_key parameters, which reduces token cost. For details, see Generic invocation.

Vector search

You can use vector_search to connect to a vector database such as Milvus and run vector search.

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

Multimodal operators

The DataFrame API lets you conveniently invoke multimodal operators to read, analyze, and transform multimodal data such as images, audio, and video. For details, see Multimodal Operators. For the API reference, see Multimodal Expressions.

The following example reads image content from URLs and performs validity checks, decoding, metadata extraction, quality scoring, and thumbnail generation.

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

Data sinks

The DataFrame write methods write the current DataFrame to the target system.

Use built-in sink functions

You can use the built-in sink functions of the DataFrame API to write data to common systems.

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

For more sink functions, see Data sinks.

Use a custom connector

You can use the write_generic function to write data through other supported connectors or custom connectors.

The connector parameter is the identifier of the connector to use, and options passes the WITH parameters of the connector.

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

If the connector requires a primary key constraint, pass primary_key:

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

Write to a Catalog table

If the target table already exists in a Catalog, you can also use write_catalog_table to write directly:

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

Write to multiple targets

When you need to write the same computation pipeline to multiple targets, use create_statement_set to create a StatementSet, pass it to the statement_set parameter of each write method, and execute them together:

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()