暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Apache Flink 2.3.0 发布:SQL 能力增强、原生 S3 与应用级生命周期管理全面升级

Flink 中文社区 2026-07-21
84
Apache Flink PMC 很兴地宣布 Apache Flink 
2.3.0
 版本发布了

Flink 
2.3.0
 实现了 15
 个 FLIP
Flink
 改进提案)的完整功能或核心功能,同时修复了多项问题并带来了众多增强改进。新版本重点提升了 SQL
 处理能力,并围绕原生 S3
 文件系统(实验性)、运行时 Rescale
、应用级生命周期管理、SQL changelog
、物化表演进与刷新控制能力增强、Metrics Exporter
 和文档重构等方向进行了系统优化。



01



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 columncomputed column以及 RENAME TO
等。这意味着当查询逻辑或表结构发生变化时,用户不再总是需要删除并重建物化表。

Flink 2.3.0
 还新增了 START_MODE
 子句,用于控制刷新作业从哪里开始处理。对于兼容的查询变化,刷新作业可以尝试从之前的 source offset
 继续执行,从而避免不必要的历史数据重处理。

  • 更多信息请参考:FLIP-550[2]FLIP-557[3]

 SinkUpsertMaterializer:显式冲突处理

SinkUpsertMaterializer
 是 Flink
 中用于处理 upsert
 流写入的重要算子。当流的 upsert key
 与 sink 
表主键不一致时,例如经过 joinprojection
 或多阶段转换后,Flink
 需要依赖该算子对数据进行协调。

此前在这类场景中,算子可能需要维护较多更新历史,带来较高状态开销,甚至存在状态膨胀风险。在 Flink 2.3.0
 中,这类执行计划默认会失败,除非用户显式通过 ON CONFLICT
 指定冲突处理策略。当前支持的策略包括 DO NOTHING
DO ERROR
 和 DO DEDUPLICATE

使用示例:

    INSERT INTO target_table
    SELECT *
    FROM source
    ON 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 events
          PARTITION BY user_id
          ORDER 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]




        02



        原生 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-2
          s3.endpoint: https://s3.us-west-2.amazonaws.com
          s3.path-style-access: false
          s3.upload.min.part.size: 5MB
          s3.upload.max.concurrent.uploads: 4
          s3.async.enabled: true
          s3.read.buffer.size: 64KB
          s3.bulk-copy.enabled: true

          此外,该实现还覆盖了 SSE-KMS
          chunked encoding
          checksum validation
           和 entropy injection
           等能力。

          • 更多信息请参考:FLIP-555[8]Native S3 FileSystem Documentation[9]




          03



          Runtime 改进与新特性

           自适应分区选择

          Flink 2.3.0
           针对 RebalancePartitioner
           和 RescalePartitioner
           引入了自适应分区选择能力。

          在数据分发过程中,Flink
           可以根据下游负载情况选择更合适的分区,从而减少流量不均和反压问题。对于存在下游处理能力波动或任务负载不均的作业,该能力有助于提升整体吞吐和稳定性。

          如需启用该能力,可参考如下配置:

            taskmanager.network.adaptive-partitioner.enabled: true
            taskmanager.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: true
                  execution.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: gzip
                      metrics.reporter.otel.batch.size: 500

                      这两个配置均为可选开启。

                      • 更多信息请参考:FLIP-553[19]




                      04



                      文档结构优化

                      Flink 2.3.0
                       对文档结构进行了重组,以提升导航和查找体验。

                      在新文档结构中,Flink SQL
                       拥有独立的一级文档入口,relational streaming
                       相关概念被提升到更显著的位置,Python
                       文档被整合到相关 API
                       区域,Contributor
                       文档也从用户文档中迁移出来。

                      已有链接会继续通过 redirect
                       保持可访问,因此用户在升级文档引用时可以平滑过渡。

                      • 更多信息请参考:FLIP-561[20]Documentation Landing Page[21]




                      05



                      升级注意事项

                      Flink
                       社区致力于确保版本升级过程尽可能顺畅。但某些变更可能需要用户在升级到 2.3.0
                       版本时,对程序、SQL
                       作业或运行时配置进行调整。

                      如果你计划从旧版本升级到 Flink 2.3.0
                      ,建议阅读官方发版说明,确认是否存在影响升级的行为变化或兼容性调整。

                      • 更多信息请参考:Flink 2.3 Release Notes[22]




                      06



                      致谢

                      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 版」 
                      复制下方链接或者扫描左边二维码
                      即可免费试用阿里云 Serverless Flink体验新一代实时计算平台的强大能力!
                      了解试用详情:https://free.aliyun.com/?productCode=sc

                      ▼ 关注「Apache Flink」 
                      回复 FFA 2026 获取大会资料

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

                      文章转载自Flink 中文社区,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                      评论