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
- Table API 概览:完整的 Table API 参考。
- 用户自定义函数:为流水线创建自定义函数。
- 流处理概念:了解动态表、时间属性等。
探索其他教程
- Flink SQL 教程:无需编写代码即可进行交互式 SQL 查询。
- DataStream API 教程:使用 DataStream API 构建有状态的流处理应用。
- Flink 运维练习场:学习操作 Flink 集群。











暂无评论内容