Apache Flink 2.3.0版本发布了。
Flink 2.3.0实现了
15个
FLIP(
Flink改进提案)的完整功能或核心功能,同时修复了多项问题并带来了众多增强改进。新版本重点提升了
SQL处理能力,并围绕原生
S3文件系统(实验性)、运行时
Rescale、应用级生命周期管理、
SQL changelog、物化表演进与刷新控制能力增强、
Metrics Exporter和文档重构等方向进行了系统优化。
Flink SQL 改进与新特性
FROM_CHANGELOG 和 TO_CHANGELOG:桥接仅追加和动态变更日志表
Flink 2.3.0
在 SQL
层新增了 changelog
转换能力,通过内置 Process Table Function
提供了与 DataStream changelog conversion API
对应的 SQL
使用接口。
FROM_CHANGELOG
可以将带有操作列的 append-only
流转换为动态表。用户可以通过 op_mapping
配置不同操作类型的映射,因此可以更方便地接入自定义 CDC
编码,或其他携带操作语义的数据流。
TO_CHANGELOG
则可以将动态表转换为 append-only changelog
流。该能力适用于审计日志、数据归档、append-only sink
,以及需要将动态表结果继续写入追加型数据管道的场景。
当前版本覆盖了基础使用场景。后续版本预计会继续补充 PARTITION BY
、invalid_op_handling
和 produces_full_deletes
等能力。
更多信息请参考:FLIP-564 [1]
物化表演进:DDL Parity 和刷新控制
物化表是 Flink SQL
中用于简化批处理和流式数据管道的一种表类型。用户只需要声明查询逻辑和刷新要求,Flink
即可自动维护相应的数据管道。
在 Flink 2.3.0
中,物化表的使用方式进一步向普通表 DDL
的使用方式靠拢。CREATE MATERIALIZED TABLE
现在支持显式定义列,同时支持声明 watermark
和 primary key
。
此外,ALTER MATERIALIZED TABLE
支持更多表结构演进操作,包括 ADD
、MODIFY
、DROP
、metadata column、computed column以及 RENAME TO
等。这意味着当查询逻辑或表结构发生变化时,用户不再总是需要删除并重建物化表。
Flink 2.3.0
还新增了 START_MODE
子句,用于控制刷新作业从哪里开始处理。对于兼容的查询变化,刷新作业可以尝试从之前的 source offset
继续执行,从而避免不必要的历史数据重处理。
更多信息请参考:FLIP-550[2]、FLIP-557[3]
SinkUpsertMaterializer:显式冲突处理
SinkUpsertMaterializer
是 Flink
中用于处理 upsert
流写入的重要算子。当流的 upsert key
与 sink
表主键不一致时,例如经过 join、projection
或多阶段转换后,Flink
需要依赖该算子对数据进行协调。
此前在这类场景中,算子可能需要维护较多更新历史,带来较高状态开销,甚至存在状态膨胀风险。在 Flink 2.3.0
中,这类执行计划默认会失败,除非用户显式通过 ON CONFLICT
指定冲突处理策略。当前支持的策略包括 DO NOTHING
、DO ERROR
和 DO DEDUPLICATE
。
使用示例:
INSERT INTO target_tableSELECT *FROM sourceON CONFLICT DO DEDUPLICATE;
更多信息请参考:FLIP-558[4]
Process Table Function 增强
Flink 2.3.0
继续增强 Process Table Function
的 SQL
处理能力。
在新版本中,迟到数据可以被处理,而不是在用户无法控制的情况下直接丢弃。同时,表参数支持使用 ORDER BY
,从而在 partition
内提供确定性的处理顺序。
使用示例:
SELECT *FROM MyTimestampedPtf(input => TABLE eventsPARTITION BY user_idORDER BY event_time);
需要注意的是,Flink 2.3.0
仅包含 FLIP-565
的部分能力,状态访问和 broadcast state
相关增强将在后续版本中继续推进。
更多信息请参考:FLIP-565[5]
UDF ARTIFACT 关键字
在 CREATE FUNCTION ... USING
中支持使用 ARTIFACT
作为 JAR
的替代写法。
使用示例:
CREATE FUNCTION my_func AS 'com.example.MyUdf'USING ARTIFACT 's3://bucket/path/my-udf.jar';
原有的 USING JAR
语法仍然有效,ARTIFACT
与 JAR
可互换使用,也可以在同一 USING
子句中按依赖来源混用。该变化让函数依赖的声明方式更加通用,也为后续更丰富的 artifact
管理能力奠定基础。
更多信息请参考:FLIP-559[6]
关键修复:MiniBatch Aggregation 记录丢失
Flink 2.3.0
修复了一个关键的 MiniBatchGroupAggFunction
问题。
在同时使用 mini-batch aggregation
和 ONE_PHASE aggregation
时,如果某个 key
对应的 batch
中只有 retract
记录,且状态中不存在该 key
,可能会导致记录丢失。该问题已在 Flink 2.3.0
中修复。
更多信息请参考:FLINK-35661[7]
原生 S3 File System
Flink 2.3.0
引入了实验性的 flink-s3-fs-native
,这是一个基于 AWS SDK v2
构建的新 S3
文件系统插件。
该插件直接使用 AWS SDK v2
,支持 IAM Roles for Service Accounts
、现代 credential provider
,以及更原生的 AWS
集成方式。同时,它采用异步、非阻塞 I/O
,并提供 FileSystem
和 RecoverableWriter
能力。
与传统 Hadoop
相关实现相比,flink-s3-fs-native
不依赖 Hadoop
,并注册了 s3://
与 s3a://
两种 scheme
。对于运行在云原生环境中的 Flink
作业,这可以降低依赖复杂度,并提升与 AWS
生态的集成体验。
配置示例:
s3.region: us-west-2s3.endpoint: https://s3.us-west-2.amazonaws.coms3.path-style-access: falses3.upload.min.part.size: 5MBs3.upload.max.concurrent.uploads: 4s3.async.enabled: trues3.read.buffer.size: 64KBs3.bulk-copy.enabled: true
此外,该实现还覆盖了 SSE-KMS
、chunked encoding
、checksum validation
和 entropy injection
等能力。
更多信息请参考:FLIP-555[8]、Native S3 FileSystem Documentation[9]
Runtime 改进与新特性
自适应分区选择
Flink 2.3.0
针对 RebalancePartitioner
和 RescalePartitioner
引入了自适应分区选择能力。
在数据分发过程中,Flink
可以根据下游负载情况选择更合适的分区,从而减少流量不均和反压问题。对于存在下游处理能力波动或任务负载不均的作业,该能力有助于提升整体吞吐和稳定性。
如需启用该能力,可参考如下配置:
taskmanager.network.adaptive-partitioner.enabled: truetaskmanager.network.adaptive-partitioner.max-traverse-size: 4
请注意,taskmanager.network.adaptive-partitioner.enabled
配置的默认值为 false
,且仅当该配置开启时taskmanager.network.adaptive-partitioner.max-traverse-size
的配置值才能生效。
更多信息请参考:FLIP-339[10]
Adaptive Scheduler 缩放历史记录与 Web UI
Flink 2.3.0
支持为开启 Adaptive Scheduler
的流作业记录并可视化扩缩容事件,其中包括并行度变化、slot
分配情况、调度器状态转换和终止原因等信息。
Web UI 新增了 Rescales 标签页,用于展示扩缩容数量、时间线、持续时间统计、单次扩缩容详情以及 Adaptive Scheduler
相关配置。通过这些信息,用户可以更直观地观察作业扩缩容过程,并分析扩缩容是否符合预期。
相关配置如下:
web.adaptive-scheduler.rescale-history.size: 100
当该配置为正整数,将启用扩缩容历史记录收集。
更多信息请参考:FLIP-487[11]、FLIP-495[12]、Elastic Scaling Rescale History[13]
面向快速回溯处理的 Watermark 对齐
Flink 2.3.0
重新设计了 watermark alignment
,以避免旧机制在处理 backlog
时造成不必要的延迟和吞吐限制。
新机制引入 watermark alignment buffer
,会将 alignment
算法的应用延后若干个更新间隔。这样在数据积压处理阶段,Flink
可以减少过早 alignment
带来的影响,从而更高效地追赶 backlog
。
需要注意的是,新的 buffer
会让 source
暂停稍晚发生,可能使窗口或 temporal
算子的状态大小轻微增加,但通常可以换来更好的 backlog
追赶能力。
相关配置如下:
pipeline.watermark-alignment.buffer-size: 3
如果将该值设置为 0
,则恢复 Flink 2.2
的行为。
更多信息请参考:FLINK-37399[14]
支持 Recovery 期间 Checkpoint
Flink 2.3.0
现在可以支持从 unaligned checkpoint
恢复期间触发 checkpoint
。
在此前版本中,checkpoint
必须等到恢复出的 channel state
被消费完成后才能完成。对于存在大量 in-flight state
的作业,这一过程可能持续较长时间,并影响恢复过程中的容错能力。
Flink 2.3.0
新增了相关配置,使作业可以在恢复期间继续进行 checkpoint
。对于拥有大量 unaligned checkpoint
状态,或者存在频繁扩缩容需求的作业,建议同时启用以下两个选项。
启用示例如下:
execution.checkpointing.unaligned.recover-output-on-downstream.enabled: trueexecution.checkpointing.unaligned.during-recovery.enabled: true
这两个配置的默认值均为 false
。其中,execution.checkpointing.unaligned.during-recovery.enabled
的配置项的开启依赖 execution.checkpointing.unaligned.recover-output-on-downstream.enabled
配置项的开启。
更多信息请参考:FLIP-547[15]
应用级生命周期管理
Flink 2.3.0
引入了 application(应用)概念,将运行时实体组织为 cluster
、application
和 job
三层结构。
cluster → application → job
Web UI
新增了 Applications
标签页,并更新了首页展示。Job
页面也会链接到其所属 application
。通过应用级生命周期管理,用户可以更清晰地理解多个 job
与同一个业务应用之间的关系。
更多信息请参考:FLIP-549[16]、FLIP-560[17]、Application Lifecycle Documentation[18]
OpenTelemetry Metrics Exporter 优化
Flink 2.3.0
提升了 OpenTelemetry exporter
在处理大体量指标 payload
时的稳定性。
新版本支持对导出的指标进行压缩和批量控制,从而减少大 payload
对指标系统带来的压力。
相关配置如下:
metrics.reporter.otel.exporter.compression: gzipmetrics.reporter.otel.batch.size: 500
这两个配置均为可选开启。
更多信息请参考:FLIP-553[19]
文档结构优化
Flink 2.3.0
对文档结构进行了重组,以提升导航和查找体验。
在新文档结构中,Flink SQL
拥有独立的一级文档入口,relational streaming
相关概念被提升到更显著的位置,Python
文档被整合到相关 API
区域,Contributor
文档也从用户文档中迁移出来。
已有链接会继续通过 redirect
保持可访问,因此用户在升级文档引用时可以平滑过渡。
更多信息请参考:FLIP-561[20]、Documentation Landing Page[21]
升级注意事项
Flink
社区致力于确保版本升级过程尽可能顺畅。但某些变更可能需要用户在升级到 2.3.0
版本时,对程序、SQL
作业或运行时配置进行调整。
如果你计划从旧版本升级到 Flink 2.3.0
,建议阅读官方发版说明,确认是否存在影响升级的行为变化或兼容性调整。
更多信息请参考:Flink 2.3 Release Notes[22]
致谢
Apache Flink
社区衷心感谢来自全球的贡献者持续的投入,让 Flink
在流批一体、云原生及 AI Native
等方向不断演进,并特别感谢所有参与 Flink 2.3.0
版本开发、评审、测试、文档和发布工作的贡献者:
Aleksandr Iushmanov, Aleksandr Savonin, AlexYinHan, Au-Miner, averyzhang, balassai,bgeng777, Biao Geng, bowenli86, bvarghese1, Cameron, Chan hae OH, David Anderson, David Radley,Dawid Wysakowicz, ddebowczyk92, Dian Fu, Dmitriy Linevich, Dong Wang, dylanhz, Efrat Levitan,Fabian Hueske, Fabian Paul, Feat Zhang, Ferenc Csaky, Gabor Somogyi, Gustavo de Morais, Hang Ruan,Hao Li, HundalTaran, ivan.torres, Jacky Lau, Jiaan Geng, Jim Hughes, Jinkun Liu, Juntao Zhang, Kunni,Leonard Xu, Liu Jiangang, Martijn Visser, Mate Czagany, Matt Cuento, mattcuento, Mika Naylor, Mingliang Liu,MUKUL GUPTA, mukul-8, Myracle, Naci Simsek, Natea Eshetu Beshada, Nflrijal, noorall, och5351, Oleksandr Nitavskyi,Pan Yuepeng, Peter Huang, Piotr Nowojski, Piotr Przybylski, Ramin Gharib, RaoraoXiong, Rion Williams, Roman,Roman Khachatryan, Rui Fan, Samrat, Samrat002, Santwana Verma, Sergey Kononov, Sergey Nuyanzin, Shekhar Prasad Rajak,Shengkai, Stepan Stepanishchev, Swapna Marru, Timo Walther, Tom Ho, Vincent-Woo, Vishal, Vishal Kamlapure, Weihua Hu,Wenjin Xie, Wren Chan, Xuannan, xuyang, Xuyang, Yi Zhang, yongliu, Yuan Mei, Yuepeng Pan,Zakelly, Zdenek Tison, Zhanghao Chen
相关链接:
[1] https://cwiki.apache.org/confluence/x/34k8G
[2] https://cwiki.apache.org/confluence/x/XwobFw
[3] https://cwiki.apache.org/confluence/x/9oPMFw
[4] https://cwiki.apache.org/confluence/x/NoTMFw
[5] https://cwiki.apache.org/confluence/x/qIo8G
[6] https://cwiki.apache.org/confluence/x/64TMFw
[7] https://issues.apache.org/jira/browse/FLINK-35661
[8] https://cwiki.apache.org/confluence/x/uYqmFw
[9] https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/deployment/filesystems/s3/
[10] https://cwiki.apache.org/confluence/x/nYyzDw
[11] https://cwiki.apache.org/confluence/x/vZCMEw
[12] https://cwiki.apache.org/confluence/x/TQr0Ew
[13] https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/deployment/elastic_scaling/#rescale-history
[14] https://issues.apache.org/jira/browse/FLINK-37399
[15] https://cwiki.apache.org/confluence/x/UQrxFg
[16] https://cwiki.apache.org/confluence/x/7gsGFw
[17] https://cwiki.apache.org/confluence/x/5AAXG
[18] https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/internals/application_lifecycle/
[19] https://cwiki.apache.org/confluence/x/1AteFw
[20] https://cwiki.apache.org/confluence/x/XII8G
[21] https://nightlies.apache.org/flink/flink-docs-release-2.3/
[22] https://nightlies.apache.org/flink/flink-docs-release-2.3/release-notes/flink-2.3/



点击「阅读原文」跳转阿里云实时计算 Flink~





