Apache Flink Table API 实时报表实战:从 Kafka 交易流到 MySQL 小时级消费报表
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Apache Flink 的 Table API 是面向批流处理的统一关系型 API同一条查询语句在无界实时流与有界批量数据上以相同语义执行并产生相同结果。本文基于 Flink 官方的 “Real Time Reporting with the Table API” 演练文档table_api.md完整拆解如何构建一个按账户跟踪财务交易的实时看板从 Kafka 读取交易数据经 Table API 聚合出小时级消费报表写入 MySQL 并通过 Grafana 所依赖的 MySQL 表。读完本文你将掌握TableEnvironment的批/流模式配置、DDL 方式注册源表与结果表、select/groupBy聚合、用户自定义函数、Tumble 窗口以及本地 Docker 一键跑通流式作业的完整流程。你将构建什么这个教程的目标是构建一个实时看板real-time dashboard按账户维度跟踪财务交易。数据流水线为Kafka(topic: transactions) → Flink Table API 聚合 → MySQL(spend_report 表) → Grafana 可视化前提条件熟悉 Java当然使用其他语言背景也能跟上熟悉SELECT、GROUP BY等基本关系型概念本地装有Java 11、Maven、Docker。{{ unstable }} 需要注意演练所使用的 Apache Flink Docker 镜像只发布了正式版本的镜像。如果你当前浏览的是 SNAPSHOT 版本文档文中涉及的版本引用将无法直接工作需要通过左侧菜单下方的 release picker 切换到最新正式版本。{{ /unstable }}获取工程演练所需的配置文件位于flink-playgrounds社区仓库中项目名flink-playground/table-walkthrough。下载后在 IDE 中打开该工程定位到核心文件SpendReport即可看到如下骨架代码EnvironmentSettings settings EnvironmentSettings.inStreamingMode(); TableEnvironment tEnv TableEnvironment.create(settings); tEnv.executeSql(CREATE TABLE transactions (\n account_id BIGINT,\n amount BIGINT,\n transaction_time TIMESTAMP(3),\n WATERMARK FOR transaction_time AS transaction_time - INTERVAL 5 SECOND\n ) WITH (\n connector kafka,\n topic transactions,\n properties.bootstrap.servers kafka:9092,\n format csv\n )); tEnv.executeSql(CREATE TABLE spend_report (\n account_id BIGINT,\n log_ts TIMESTAMP(3),\n amount BIGINT\n, PRIMARY KEY (account_id, log_ts) NOT ENFORCED ) WITH (\n connector jdbc,\n url jdbc:mysql://mysql:3306/sql-demo,\n table-name spend_report,\n driver com.mysql.jdbc.Driver,\n username sql-demo,\n password demo-sql\n )); Table transactions tEnv.from(transactions); report(transactions).executeInsert(spend_report);其中report方法此时尚未实现它正是我们要填充业务逻辑的位置。提示如果在 Windows 上运行 Docker 且数据生成容器table-walkthrough_data-generator_1启动失败请确认使用了正确的 shell。例如docker-entrypoint.sh需要 bash若环境中不可用会抛出standard_init_linux.go:211: exec user process caused no such file or directory错误。变通办法是把docker-entrypoint.sh第一行的 shell 切换为sh。拆解代码执行环境前两行创建了TableEnvironment。它是作业的属性配置入口可以指定批或流执行模式、创建各类源。本演练创建的是标准的流式执行环境EnvironmentSettings settings EnvironmentSettings.inStreamingMode(); TableEnvironment tEnv TableEnvironment.create(settings);结合源码可以看到这两种快捷方式的本质。在 EnvironmentSettings.java 中public static EnvironmentSettings inStreamingMode() { return EnvironmentSettings.newInstance().inStreamingMode().build(); } public static EnvironmentSettings inBatchMode() { return EnvironmentSettings.newInstance().inBatchMode().build(); }两者都是Builder的快捷方式最终分别把ExecutionOptions.RUNTIME_MODE配置项设置为STREAMING或BATCH。从 Javadoc 的注释可以看出两种模式的边界流模式inStreamingMode既可以处理有界数据流也可以处理无界数据流批模式inBatchMode专为批场景优化只能处理有界数据流。此外Builder还支持withBuiltInCatalogName(...)、withBuiltInDatabaseName(...)、withConfiguration(...)、withClassLoader(...)等高级配置见 EnvironmentSettings.java分别用于指定初始 Catalog 名、初始数据库名、追加额外配置、以及代码生成与 UDF 加载使用的 ClassLoader。这些参数只在实例化TableEnvironment时生效之后不可更改。拆解代码注册表Catalog 与 DDL接下来表被注册到当前 Catalog 中。Table Source 提供对数据库、键值存储、消息队列、文件系统等外部系统数据的读取能力Table Sink 则把表写出到外部存储系统。根据 Source 与 Sink 的类型不同它们支持 CSV、JSON、Avro、Parquet 等不同格式。输入表Kafka 上的交易数据tEnv.executeSql(CREATE TABLE transactions (\n account_id BIGINT,\n amount BIGINT,\n transaction_time TIMESTAMP(3),\n WATERMARK FOR transaction_time AS transaction_time - INTERVAL 5 SECOND\n ) WITH (\n connector kafka,\n topic transactions,\n properties.bootstrap.servers kafka:9092,\n format csv\n ));transactions表让我们读取信用卡交易数据包含账户 IDaccount_id、时间戳transaction_time和美元金额amount。它本质上是对一个名为transactions的 Kafka topic 的逻辑视图数据为 CSV 格式。几个关键 DDL 要素要素作用account_id BIGINT/amount BIGINT业务字段账户 ID 与金额transaction_time TIMESTAMP(3)毫秒精度的事件时间列WATERMARK FOR transaction_time AS transaction_time - INTERVAL 5 SECOND事件时间水位线容忍 5 秒数据乱序是后续窗口聚合能正确触发的前提connector kafka使用 Kafka 连接器topic transactions订阅的 topicproperties.bootstrap.servers kafka:9092Kafka bootstrap 地址Docker 网络内的服务名format csv消息的序列化格式输出表MySQL 中的报表数据tEnv.executeSql(CREATE TABLE spend_report (\n account_id BIGINT,\n log_ts TIMESTAMP(3),\n amount BIGINT\n, PRIMARY KEY (account_id, log_ts) NOT ENFORCED ) WITH (\n connector jdbc,\n url jdbc:mysql://mysql:3306/sql-demo,\n table-name spend_report,\n driver com.mysql.jdbc.Driver,\n username sql-demo,\n password demo-sql\n ));第二张表spend_report存放聚合的最终结果其底层存储是 MySQL 中的一张表。声明PRIMARY KEY (account_id, log_ts) NOT ENFORCED的原因在于聚合结果在流上是不断更新的JDBC Sink 需要主键信息才能以“更新”而非纯追加语义写入数据库避免同一账户同一小时的结果重复累积。查询骨架from 与 executeInsert环境配置好、表注册完成后就可以搭建第一个应用从TableEnvironment读取输入表再通过executeInsert把结果写入输出表。Table transactions tEnv.from(transactions); report(transactions).executeInsert(spend_report);在 API 层面executeInsert是Table接口的默认方法见 Table.java除executeInsert(String tablePath)外还提供了带overwrite参数及以TableDescriptor为目标的重载。对于流式作业executeInsert会启动一个持续消费输入、持续产出结果的作业对开发期调试注释中也提示可以用toAppendStream/toChangelogStream等方式直接查看中间结果。批模式开发与测试Flink 的一个独特属性是批与流之间的一致语义你可以在静态数据集上用批模式开发和测试应用然后以流式应用的形式部署到生产环境。演练工程自带一个测试类SpendReportTest用来验证报表逻辑它创建的是批模式的表环境EnvironmentSettings settings EnvironmentSettings.inBatchMode(); TableEnvironment tEnv TableEnvironment.create(settings);这意味着后续每次修改report方法都可以先在批模式测试中验证正确性而不用启动整套 Kafka/MySQL 环境。第一次实现小时级消费报表现在骨架就绪可以添加业务逻辑了。目标是展示每个账户在一天中每个小时内的总消费额。因此时间戳列需要从毫秒粒度取整向下取整到小时粒度。Flink 支持用纯 SQL 或 Table API 开发关系型应用。Table API 是一个受 SQL 启发、可用 Java 或 Python 编写的流畅 DSL支持良好的 IDE 集成。和 SQL 查询一样Table 程序可以选择所需字段、按键分组配合 内建函数 如floor和sum即可写出这份报表public static Table report(Table transactions) { return transactions.select( $(account_id), $(transaction_time).floor(TimeIntervalUnit.HOUR).as(log_ts), $(amount)) .groupBy($(account_id), $(log_ts)) .select( $(account_id), $(log_ts), $(amount).sum().as(amount)); }逻辑分三步select中用floor(TimeIntervalUnit.HOUR)把transaction_time向下取整到小时并别名为log_tsgroupBy(account_id, log_ts)按“账户 小时”分组最后一次select对amount做sum聚合。表达式floor的实现在 BaseExpressions.java 中其注释给出示例lit(12:44:31).toDate().floor(MINUTE)得到12:44:00——时间取整的语义与 SQL 中的时间函数一致。用该实现运行批模式测试测试通过。用户自定义函数UDFFlink 内置函数数量有限有时需要以 用户自定义函数 扩展。假如floor不是预定义的你可以自行实现一个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( $(account_id), call(MyFloor.class, $(transaction_time)).as(log_ts), $(amount)) .groupBy($(account_id), $(log_ts)) .select( $(account_id), $(log_ts), $(amount).sum().as(amount)); }这个查询会消费transactions表的全部记录计算报表并以高效、可扩展的方式输出结果。运行测试该实现同样通过。加入窗口Tumble 窗口按时间对数据分组是数据处理中的典型操作尤其在处理无限流时。基于时间的分组称为窗口Flink 提供了灵活的窗口语义。最基础的窗口类型是Tumble滚动窗口具有固定大小且窗口桶之间不重叠。public static Table report(Table transactions) { return transactions .window(Tumble.over(lit(1).hour()).on($(transaction_time)).as(log_ts)) .groupBy($(account_id), $(log_ts)) .select( $(account_id), $(log_ts).start().as(log_ts), $(amount).sum().as(amount)); }这定义了应用基于时间戳列使用一小时滚动窗口。于是时间戳为2019-06-01 01:23:47的行会落入2019-06-01 01:00:00对应的窗口输出时$(log_ts).start()取窗口的起始时间作为报表的时间标签。时间维度上的聚合有其独特性与floor或你的 UDF 不同窗口函数是内建intrinsic算子运行时可以据此施加额外优化。Table.window(GroupWindow)的 Javadoc见 Table.java也明确了批流两种语境下的含义流表无限表将记录分组到时间或行区间定义的窗口中是“定义有限聚合”的必要手段且只有当groupBy(...)中除窗口别名外还包含额外的分组属性时窗口聚合才是并行操作否则整条流会被单并发处理批表有限表窗口化本质上等价于基于时间属性的groupBy的便捷 API。用该实现运行测试测试依然通过。流模式运行Docker 一键启动到此为止一个功能完备、有状态、分布式、流式的应用就完成了。查询会持续消费 Kafka 上的交易流计算每小时消费额并在结果就绪时立即发出。由于输入是无界的查询会一直运行直到被手动停止。又因为作业使用了基于时间窗口的聚合Flink 可以执行特定优化——例如当框架知道某个窗口不会再有记录到来时就会清理对应的状态。演练工程完全 Docker 化可以在本地作为流式应用运行。环境包含一个 Kafka topic、一个持续的数据生成器、MySQL以及 Grafana。在table-walkthrough目录下启动 docker-compose$ docker-compose build $ docker-compose up -d运行中的作业信息可以通过 Flink Web 控制台http://localhost:8082/查看。接着可以直接进入 MySQL 验证结果$ docker-compose exec mysql mysql -Dsql-demo -usql-demo -pdemo-sql mysql use sql-demo; Database changed mysql select count(*) from spend_report; ---------- | count(*) | ---------- | 110 | ----------最后打开 Grafanahttp://localhost:3000/d/FOe0PbmGk/walkthrough?viewPanel2orgId1refresh5s即可看到完全可视化后的报表效果。小结本演练完整呈现了 Table API 开发实时报表的标准路径用EnvironmentSettings.inStreamingMode()/inBatchMode()区分流式生产作业与批式测试环境两者共享同一套查询代码与语义用executeSql以 DDL 注册 Kafka 源表含 watermark与 JDBC 结果表含主键以支持更新语义用selectfloor/groupBysum写出小时级聚合必要时以ScalarFunction自定义扩展函数用Tumble滚动窗口替代手工时间取整获得运行时可优化的内建窗口算子并让 Flink 在窗口结束后自动清理状态通过executeInsert提交作业用 Docker 化的 Kafka MySQL Grafana 环境端到端验证。这套“批模式开发测试、流模式生产部署、统一关系语义”的工作方式正是 Flink Table API 用于数据分析、数据管道与 ETL 场景的核心价值所在。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Table API 实时报表实战从 Kafka 到 MySQL 再到 Grafana 的端到端看板搭建Flink Table API 实时报表实战从 Kafka 到 MySQL 再到 Grafana 的端到端看板搭建 Apache Flink 的 Table大数据流处理批处理数据工程Apache Doris数据报表实时报表生成与展示Apache Doris数据报表实时报表生成与展示 你是否还在为传统报表系统的延迟问题烦恼当业务需要实时决策时T1的报表数据早已失去价值。Apache数据库OLAP大数据数据仓库分布式数据库实时分析列式数据库使用 Flink CDC 构建实时数据湖MySQL 分库分表到 Apache Iceberg 的实战教程使用 Flink CDC 构建实时数据湖MySQL 分库分表到 Apache Iceberg 的实战教程 导读 本教程基于当前开源仓库 flink cdcF后端数据集成大数据流处理变更数据捕获数据同步上一篇k-skill 之 daangn-cars-search基于 당근중고차 Remix _data 路由的只读二手车检索实现下一篇重塑数字主权Win11Debloat如何重新定义Windows体验治理范式创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考