Table API 教程

Apache Flink 提供 Table API,作为统一的关系型 API,用于批处理和流处理。无论面对无界的实时数据流,还是有界的批数据集,查询都使用相同的语义执行,并产生相同的结果。Flink 的 Table API 常用于数据分析、数据流水线和 ETL 应用。

你将构建什么

在本教程中,你将构建一份消费报表,按账户和小时汇总交易金额:

Transactions (generated data) → Flink (Table API aggregation) → Console (results)

你将学会:

  • 设置 Table API 流处理环境。
  • 使用 Table API 创建表。
  • 通过 Table API 操作编写持续聚合。
  • 实现用户自定义函数(UDF)。
  • 使用基于时间的窗口进行聚合。
  • 在批处理模式下测试流处理应用。

前提条件

本教程假设你对 Java 或 Python 有一定了解;即使你使用其他编程语言,也应能跟上教程。你还需要熟悉基本的关系型概念,例如 SELECT 和 GROUP BY 子句。

遇到困难怎么办

如果遇到困难,可以查看社区支持资源。尤其是 Apache Flink 的用户邮件列表,它一直是 Apache 项目中最活跃的邮件列表之一,是迅速获得帮助的好途径。

如何跟着教程操作

如果要跟着操作,你需要一台具备以下环境的计算机:

Java

  • Java 11、17 或 21。
  • Maven。

Python

  • Java 11、17 或 21。
  • Python 3.9、3.10、3.11 或 3.12。

Java:创建项目

Flink 提供的 Maven Archetype 可以迅速生成包含全部必要依赖的项目骨架,让你专注于填写业务逻辑。

$ mvn archetype:generate \
    -DarchetypeGroupId=org.apache.flink \
    -DarchetypeArtifactId=flink-walkthrough-table-java \
    -DarchetypeVersion=2.3.0 \
    -DgroupId=spendreport \
    -DartifactId=spendreport \
    -Dversion=0.1 \
    -Dpackage=spendreport \
    -DinteractiveMode=false

你可以按需修改 groupId、artifactId 和 package。使用上述参数,Maven 会创建名为 spendreport 的文件夹,其中的项目包含完成本教程所需的全部依赖。

把项目导入编辑器后,你可以找到 SpendReport.java 文件,其代码如下,可直接在 IDE 内运行。

在 IDE 中运行:如果遇到 java.lang.NoClassDefFoundError 异常,很可能是 classpath 中缺少必要的 Flink 依赖。

  • IntelliJ IDEA:进入 Run > Edit Configurations > Modify options,勾选 include dependencies with “Provided” scope。

Python:安装 PyFlink

使用 Python Table API 需要安装 PyFlink。它发布在 PyPI 上,可以方便地通过 pip 安装:

$ python -m pip install apache-flink

提示:建议在虚拟环境中安装 PyFlink,以隔离项目依赖。

安装完成后,新建 spend_report.py 文件,用于编写 Table API 程序。

完整程序

消费报表程序的完整代码如下:

Java

public class SpendReport {

    public static void main(String[] args) throws Exception {
        // Create a Table environment for streaming
        EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
        TableEnvironment tEnv = TableEnvironment.create(settings);

        // Create the source table
        // The DataGen connector generates an infinite stream of transactions
        tEnv.createTemporaryTable("transactions",
            TableDescriptor.forConnector("datagen")
                .schema(Schema.newBuilder()
                    .column("accountId", DataTypes.BIGINT())
                    .column("amount", DataTypes.BIGINT())
                    .column("transactionTime", DataTypes.TIMESTAMP(3))
                    .watermark("transactionTime", "transactionTime - INTERVAL '5' SECOND")
                    .build())
                .option("rows-per-second", "100")
                .option("fields.accountId.min", "1")
                .option("fields.accountId.max", "5")
                .option("fields.amount.min", "1")
                .option("fields.amount.max", "1000")
                .build());

        // Read from the source table
        Table transactions = tEnv.from("transactions");

        // Apply the business logic
        Table result = report(transactions);

        // Print the results to the console
        result.execute().print();
    }

    public static Table report(Table transactions) {
        return transactions
            .select(
                $("accountId"),
                $("transactionTime").floor(TimeIntervalUnit.HOUR).as("logTs"),
                $("amount"))
            .groupBy($("accountId"), $("logTs"))
            .select(
                $("accountId"),
                $("logTs"),
                $("amount").sum().as("amount"));
    }
}

Python

from pyflink.table import TableEnvironment, EnvironmentSettings, TableDescriptor, Schema, DataTypes
from pyflink.table.expression import TimeIntervalUnit
from pyflink.table.expressions import col

def main():
    # Create a Table environment for streaming
    settings = EnvironmentSettings.in_streaming_mode()
    t_env = TableEnvironment.create(settings)

    # Write all results to one file for easier viewing
    t_env.get_config().set("parallelism.default", "1")

    # Create the source table
    # The DataGen connector generates an infinite stream of transactions
    t_env.create_temporary_table(
        "transactions",
        TableDescriptor.for_connector("datagen")
            .schema(Schema.new_builder()
                .column("accountId", DataTypes.BIGINT())
                .column("amount", DataTypes.BIGINT())
                .column("transactionTime", DataTypes.TIMESTAMP(3))
                .watermark("transactionTime", "transactionTime - INTERVAL '5' SECOND")
                .build())
            .option("rows-per-second", "100")
            .option("fields.accountId.min", "1")
            .option("fields.accountId.max", "5")
            .option("fields.amount.min", "1")
            .option("fields.amount.max", "1000")
            .build())

    # Read from the source table
    transactions = t_env.from_path("transactions")

    # Apply the business logic
    result = report(transactions)

    # Print the results to the console
    result.execute().print()

def report(transactions):
    return transactions \
        .select(
            col("accountId"),
            col("transactionTime").floor(TimeIntervalUnit.HOUR).alias("logTs"),
            col("amount")) \
        .group_by(col("accountId"), col("logTs")) \
        .select(
            col("accountId"),
            col("logTs"),
            col("amount").sum.alias("amount"))

if __name__ == '__main__':
    main()

逐步解读代码

执行环境

最初几行设置 TableEnvironment。你可以通过表环境设置作业属性,指定编写的是批处理还是流处理应用,并创建数据源。本教程创建一个采用流式执行的标准表环境。

Java

EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
TableEnvironment tEnv = TableEnvironment.create(settings);

Python

settings = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(settings)

创建表

接下来,创建一张表示交易数据的表。DataGen 连接器会生成无限的随机交易数据流。

Java

使用 TableDescriptor,可以通过代码定义表:

tEnv.createTemporaryTable("transactions",
    TableDescriptor.forConnector("datagen")
        .schema(Schema.newBuilder()
            .column("accountId", DataTypes.BIGINT())
            .column("amount", DataTypes.BIGINT())
            .column("transactionTime", DataTypes.TIMESTAMP(3))
            .watermark("transactionTime", "transactionTime - INTERVAL '5' SECOND")
            .build())
        .option("rows-per-second", "100")
        .option("fields.accountId.min", "1")
        .option("fields.accountId.max", "5")
        .option("fields.amount.min", "1")
        .option("fields.amount.max", "1000")
        .build());

Python

使用 TableDescriptor,可以通过代码定义表:

t_env.create_temporary_table(
    "transactions",
    TableDescriptor.for_connector("datagen")
        .schema(Schema.new_builder()
            .column("accountId", DataTypes.BIGINT())
            .column("amount", DataTypes.BIGINT())
            .column("transactionTime", DataTypes.TIMESTAMP(3))
            .watermark("transactionTime", "transactionTime - INTERVAL '5' SECOND")
            .build())
        .option("rows-per-second", "100")
        .option("fields.accountId.min", "1")
        .option("fields.accountId.max", "5")
        .option("fields.amount.min", "1")
        .option("fields.amount.max", "1000")
        .build())

transactions 表生成的信用卡交易具有以下字段:

  • accountId:1 到 5 之间的账户 ID。
  • amount:1 到 1000 之间的交易金额。
  • transactionTime:带有水位线的时间戳,用于处理迟到数据。

查询

配置环境并注册表后,就可以构建第一个应用了。你可以从 TableEnvironment 读取输入表,并应用 Table API 操作。业务逻辑在 report 函数中实现。

Java

Table transactions = tEnv.from("transactions");
Table result = report(transactions);
result.execute().print();

Python

transactions = t_env.from_path("transactions")
result = report(transactions)
result.execute().print()

实现报表

作业骨架搭好后,就可以添加业务逻辑。目标是构建一份报表,展示一天中每个小时、每个账户的总消费金额。这意味着需要把时间戳列从毫秒粒度向下取整到小时粒度。

Flink 支持使用纯 SQL 或 Table API 开发关系型应用。Table API 是受 SQL 启发的流式 DSL,可用 Java 或 Python 编写,并与 IDE 深度集成。与 SQL 查询一样,Table 程序可以选择所需字段,并按键分组。这些功能结合 floor、sum 等内置函数,就能实现这份报表。

Java

public static Table report(Table transactions) {
    return transactions
        .select(
            $("accountId"),
            $("transactionTime").floor(TimeIntervalUnit.HOUR).as("logTs"),
            $("amount"))
        .groupBy($("accountId"), $("logTs"))
        .select(
            $("accountId"),
            $("logTs"),
            $("amount").sum().as("amount"));
}

Python

def report(transactions):
    return transactions \
        .select(
            col("accountId"),
            col("transactionTime").floor(TimeIntervalUnit.HOUR).alias("logTs"),
            col("amount")) \
        .group_by(col("accountId"), col("logTs")) \
        .select(
            col("accountId"),
            col("logTs"),
            col("amount").sum.alias("amount"))

测试

Java

项目包含测试类 SpendReportTest,它在批处理模式下使用静态数据验证报表逻辑。

EnvironmentSettings settings = EnvironmentSettings.inBatchMode();
TableEnvironment tEnv = TableEnvironment.create(settings);

// Create test data using fromValues
Table transactions = tEnv.fromValues(
    DataTypes.ROW(
        DataTypes.FIELD("accountId", DataTypes.BIGINT()),
        DataTypes.FIELD("amount", DataTypes.BIGINT()),
        DataTypes.FIELD("transactionTime", DataTypes.TIMESTAMP(3))
    ),
    Row.of(1L, 188L, LocalDateTime.of(2024, 1, 1, 9, 0, 0)),
    Row.of(1L, 374L, LocalDateTime.of(2024, 1, 1, 9, 30, 0)),
    // ... more test data
);

Python

你可以切换到批处理模式,并使用静态测试数据来测试 report 函数。创建单独的测试文件,例如 test_spend_report.py:

from datetime import datetime
from pyflink.table import TableEnvironment, EnvironmentSettings, DataTypes

from spend_report import report

def test_report():
    settings = EnvironmentSettings.in_batch_mode()
    t_env = TableEnvironment.create(settings)

    # Create test data using from_elements
    transactions = t_env.from_elements(
        [
            (1, 188, datetime(2024, 1, 1, 9, 0, 0)),
            (1, 374, datetime(2024, 1, 1, 9, 30, 0)),
            (2, 200, datetime(2024, 1, 1, 9, 15, 0)),
        ],
        DataTypes.ROW([
            DataTypes.FIELD("accountId", DataTypes.BIGINT()),
            DataTypes.FIELD("amount", DataTypes.BIGINT()),
            DataTypes.FIELD("transactionTime", DataTypes.TIMESTAMP(3))
        ])
    )

    # Test the report function
    result = report(transactions)

    # Collect results and verify
    rows = [row for row in result.execute().collect()]
    assert len(rows) == 2  # Two accounts

if __name__ == '__main__':
    test_report()
    print("All tests passed!")

使用 python test_spend_report.py 或 pytest test_spend_report.py 运行。

Flink 的一个独特之处在于,它在批处理和流处理之间提供一致的语义。因此,你可以在批处理模式下用静态数据集开发和测试应用,再作为流处理应用部署到生产环境。

用户自定义函数

Flink 提供了许多内置函数,但有时你需要通过用户自定义函数扩展它。如果没有预定义的 floor,你也可以自己实现。

Java

import java.time.LocalDateTime;
import java.time.temporal.ChronoUnit;

import org.apache.flink.table.annotation.DataTypeHint;
import org.apache.flink.table.functions.ScalarFunction;

public class MyFloor extends ScalarFunction {

    public @DataTypeHint("TIMESTAMP(3)") LocalDateTime eval(
        @DataTypeHint("TIMESTAMP(3)") LocalDateTime timestamp) {

        return timestamp.truncatedTo(ChronoUnit.HOURS);
    }
}

然后迅速将它集成到应用中:

public static Table report(Table transactions) {
    return transactions
        .select(
            $("accountId"),
            call(MyFloor.class, $("transactionTime")).as("logTs"),
            $("amount"))
        .groupBy($("accountId"), $("logTs"))
        .select(
            $("accountId"),
            $("logTs"),
            $("amount").sum().as("amount"));
}

Python

from pyflink.table.udf import udf
from datetime import datetime

@udf(result_type=DataTypes.TIMESTAMP(3))
def my_floor(timestamp: datetime) -> datetime:
    return timestamp.replace(minute=0, second=0, microsecond=0)

然后将它集成到应用中:

def report(transactions):
    return transactions \
        .select(
            col("accountId"),
            my_floor(col("transactionTime")).alias("logTs"),
            col("amount")) \
        .group_by(col("accountId"), col("logTs")) \
        .select(
            col("accountId"),
            col("logTs"),
            col("amount").sum.alias("amount"))

该查询会消费 transactions 表中的全部记录,计算报表,并以高效、可扩展的方式输出结果。使用这个实现运行测试,将会通过。

处理表函数(仅 Java)

对于更高级的逐行处理,Flink 提供处理表函数(Process Table Functions,PTF)。PTF 可以转换表中的每一行,并能使用状态和定时器等强大功能。下面是一个简单的无状态示例,用于筛选高金额交易并将其格式化:

import org.apache.flink.table.annotation.ArgumentHint;
import org.apache.flink.table.annotation.ArgumentTrait;
import org.apache.flink.table.functions.ProcessTableFunction;
import org.apache.flink.types.Row;

public class HighValueAlerts extends ProcessTableFunction<String> {

    private static final long HIGH_VALUE_THRESHOLD = 500;

    public void eval(@ArgumentHint(ArgumentTrait.ROW_SEMANTIC_TABLE) Row transaction) {
        Long amount = transaction.getFieldAs("amount");
        if (amount > HIGH_VALUE_THRESHOLD) {
            Long accountId = transaction.getFieldAs("accountId");
            collect("Alert: Account " + accountId + " made a high-value transaction of " + amount);
        }
    }
}

你可以修改 main() 方法来尝试这个 PTF。替换现有的 report() 调用,或在旁边增加以下代码:

// Instead of (or in addition to) the aggregation report:
// Table result = report(transactions);
// result.execute().print();

// Try the PTF to see high-value transaction alerts:
Table alerts = transactions.process(HighValueAlerts.class);
alerts.execute().print();

这段代码只会针对金额超过 500 的交易输出告警。PTF 结合状态和定时器后,能进一步实现复杂的事件驱动逻辑。更多高级示例见 PTF 文档。

注意:处理表函数目前仅支持 Java。Python 可以使用用户自定义表函数(UDTF)实现类似的逐行处理。

添加窗口

按时间对数据分组,是数据处理中常见的操作,处理无限数据流时尤其如此。基于时间的分组称为窗口,Flink 提供灵活的窗口语义。最基础的一种窗口称为滚动窗口(Tumble),其大小固定,各窗口之间不重叠。

尝试修改 report() 函数,用窗口代替 floor():

Java

public static Table report(Table transactions) {
    return transactions
        .window(Tumble.over(lit(10).seconds()).on($("transactionTime")).as("logTs"))
        .groupBy($("accountId"), $("logTs"))
        .select(
            $("accountId"),
            $("logTs").start().as("logTs"),
            $("amount").sum().as("amount"));
}

Python

from pyflink.table.expressions import col, lit
from pyflink.table.window import Tumble

def report(transactions):
    return transactions \
        .window(Tumble.over(lit(10).seconds).on(col("transactionTime")).alias("logTs")) \
        .group_by(col("accountId"), col("logTs")) \
        .select(
            col("accountId"),
            col("logTs").start.alias("logTs"),
            col("amount").sum.alias("amount"))

这将应用定义为使用基于时间戳列的 10 秒滚动窗口。因此,时间戳为 2024-01-01 01:23:47 的一行,会进入 2024-01-01 01:23:40 开始的窗口。

基于时间的聚合有其特殊性:与其他属性不同,在持续运行的流处理应用中,时间通常会向前推进。与 floor 和你的 UDF 不同,窗口函数是运行时内建机制,允许运行时应用额外的优化。在批处理场景中,窗口提供了方便的 API,可按时间戳属性对记录分组。

修改后运行应用,就能看到每 10 秒输出一次的窗口结果。

运行应用

至此,一个功能完整、具有状态的分布式流处理应用就完成了!查询会持续生成交易,计算按小时汇总的消费金额,并在结果准备好时立即输出。由于输入无界,查询会一直运行,直到手动停止。

Java

在 IDE 中运行 SpendReport 类,就能看到控制台打印的流处理结果。

Python

从命令行运行程序:

$ python spend_report.py

该命令会构建 Python Table API 程序,并在本地 mini cluster 中运行。你也可以把 Python Table API 程序提交到远程集群。详情参见作业提交示例。

后续学习

恭喜你完成本教程!下面是继续学习的一些方向:

进一步了解 Table API

探索其他教程

生产部署

  • 部署概览:在生产环境中部署 Flink。
  • 连接器:连接 Kafka、数据库、文件系统等。

原文:Apache Flink 文档贡献者,Table API Tutorial(release-2.3)。本文件对原文正文作中文翻译,代码保持原样;2026-10-03。原文版权与署名保留,按 Apache License 2.0 使用。许可副本见随附 LICENSE.txt;与所复制内容相关的 NOTICE 如下:

Apache Flink
Copyright 2014-2026 The Apache Software Foundation
This product includes software developed at
The Apache Software Foundation (http://www.apache.org/).

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容