有赞实时计算 Flink 1.13 升级实践

**作者:** LiChuang **来源:** 有赞技术团队,2021 年 12 月 3 日 **原文:** https://tech.youzan.com/flink_13/

本文整理的是有赞将实时计算引擎从 Flink 1.10 升级到 1.13.2 的历史实践,针对其 SQL 平台、自研连接器和集群环境。它不是适用于所有 Flink 集群的通用升级脚本;文中的版本结论、兼容性行为和操作应结合当前官方文档与实际环境重新确认。

为什么升级

有赞的实时任务逐步转为 Flink SQL 后,Flink 1.10 的 SQL 能力已无法满足部分业务需求。原有任务运行在 YARN 上,公司推进应用容器化后,希望让实时 SQL 作业进入统一的 Kubernetes 资源池,以便隔离任务并弹性调度。团队也希望借升级采用社区持续完善的 Flink on Kubernetes 能力,因此决定把引擎从 1.10 升至 1.13.2。

这次升级关注的能力

SQL 开发和时间语义

Flink 1.10 之后,SQL connector 属性的写法有所简化。Flink 1.13 还调整了时间函数的时区语义,并支持 TIMESTAMP_LTZ:CURRENT_TIMESTAMP、CURRENT_TIME、CURRENT_DATE、NOW() 和 PROCTIME() 会按本地时区处理;例如 CURRENT_TIMESTAMP 返回 TIMESTAMP_LTZ,而不是旧版本中的 TIMESTAMP。在窗口处理中结合 TIMESTAMP 与 TIMESTAMP_LTZ,也有助于处理夏令时。

窗口表值函数(Window TVF)使窗口起止时间可以参与后续查询,因此除了普通窗口聚合和 Join,也能按窗口做 Top-K。原文给出的例子如下:


SELECT *
FROM (
  SELECT *, ROW_NUMBER() OVER (
    PARTITION BY window_start, window_end ORDER BY price DESC
  ) AS rownum
  FROM (
    SELECT window_start, window_end, supplier_id,
           SUM(price) AS price, COUNT(*) AS cnt
    FROM TABLE(
      TUMBLE(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '10' MINUTES)
    )
    GROUP BY window_start, window_end, supplier_id
  )
) WHERE rownum <= 3;

文章还介绍了累积窗口 CUMULATE。它按固定步长扩展统计区间,直到最大窗口长度;后续窗口可以纳入此前窗口迟到的数据,并复用已有聚合结果:


SELECT window_time, window_start, window_end, SUM(price) AS total_price
FROM TABLE(
  CUMULATE(TABLE Bid, DESCRIPTOR(bidtime),
           INTERVAL '2' MINUTES, INTERVAL '10' MINUTES)
)
GROUP BY window_start, window_end, window_time;

在这个例子中,第一个窗口统计第一段区间,第二个窗口合并前两段,第三个窗口合并前三段。窗口步长为 2 分钟,最大窗口为 10 分钟。

Hive、Kubernetes 与状态管理

团队也看重 Flink on Hive 的生产能力:Hive 方言可支持常见 DML、DQL 以及 Hive 脚本迁移;FileSystem connector 的改进使 Table/SQL sink 能写入 CSV、JSON、Avro、Parquet、ORC 等格式及 Hive 表格式。原文将这些能力作为此次版本升级带来的收益。

原文列出的 Flink on Kubernetes 进展包括基于 Kubernetes 内置能力的 JobManager 高可用方案,以及 Application 模式。Application 模式按应用启动集群,作业图生成和提交移至 JobManager 端执行,从而分散客户端侧的网络上传下载负载。作者也希望实时任务迁入 K8s 后,能与离线任务共同弹性伸缩,改善资源利用率。

在状态管理方面,文章提到 Flink 1.10 之后 Savepoint 中元数据和状态数据使用相对路径,便于迁移目录;Flink 1.12 的 Unaligned Checkpoint 可在使用保留检查点时扩缩容,作者认为它有助于减轻反压下检查点超时。Flink 1.13 还可报告失败或取消的 Checkpoint 统计,便于定位失败原因。

Upsert Kafka 与其他运维能力

Upsert Kafka 连接器可以消费变更日志:作为 source 时,同一 key 的 value 表示该 key 的最新值,存在 key 时视作更新,不存在时视作插入;作为 sink 时,插入和更新写为普通 Kafka 消息,删除写为 value 为空的消息。Flink 按主键分区,使同一主键的更新和删除进入同一分区并保持有序。作者将它视为处理部分重复数据场景的一种能力,但它本身不能替代端到端幂等设计。

升级还带来了更多 Source/Sink 格式支持,例如 Kafka Raw 格式和文件系统调试时使用 JSON。Web UI 可展示 JobManager、TaskManager 的内存指标;背压图使用更直观的颜色和繁忙度、背压比例,另有 CPU 火焰图以及作业历史异常列表。

升级过程:先改造连接器,再迁移 SQL

自定义连接器与 UDF

有赞的数据链路包含 NSQ、Kafka、MySQL、TiDB、ClickHouse、HBase 等组件。此次升级需要维护多种官方未提供或经过定制的连接器,包括 NSQ connector、无用户名密码的 JDBC connector、ClickHouse connector 和高可用 HBase connector。

Flink 1.10 中,DataStream 与 Table connector 使用 Row 数据结构;从 Flink 1.11 开始,FLIP-95 重构 TableSource/TableSink API 并引入 SQL 内部数据结构 RowData。因此,团队需要重构上述定制的 Table connector。无凭据 JDBC connector 使用连接池时还遇到一个实际问题:若任务长时间没有数据,连接可能被释放;数据恢复后继续用旧连接写入会报连接已关闭,造成任务反复重启。作者的处理办法是在 flush 前检查连接有效性,失效时重新建连。原文以截图展示了相关代码,未提供可复制的完整代码段,因此这里不补写实现。

对于 UDF,团队主要把 Maven 中的 flink-table-common 依赖调整到对应 Flink 版本。Flink 1.13 中,若 UDF 参数是 Object,需增加类型提示:


@DataTypeHint(inputGroup = InputGroup.ANY)

批量转换 SQL 时重点核对三类变化

团队为数百个既有 Flink 1.10 SQL 任务编写转换流程,自动生成 Flink 1.13 语法,再批量进行语法检查。转换重点包括 connector 配置项、时间函数及其类型、Upsert sink 的主键。

第一类是 connector 属性名。例如 Kafka 配置中,format.fail-on-not-json-record = false 对应新写法 json.ignore-parse-errors = true;connector.startup-mode = earliest-offset 对应 scan.startup.mode = earliest-offset。含义相近并不代表配置 key 可以原样保留,迁移时要逐项映射。

第二类是时间函数的语义和类型。原文称,旧任务可能曾给 CURRENT_TIMESTAMP 结果加 8 小时或减 16 小时,以修正旧版本按 UTC+0 处理带来的时差;升级后需要识别这些任务,再按新时区语义调整。不能对所有 SQL 机械增减八小时,应按作业时区和业务时间口径逐条确认。

类型上,Flink 1.13 的 CURRENT_TIMESTAMP 返回 TIMESTAMP_LTZ。与 TIMESTAMP 混用时,TIMESTAMPDIFF、TIMESTAMPADD 等函数可能报类型不匹配。原文举例:


TIMESTAMPDIFF(
  MINUTE,
  (current_timestamp - INTERVAL '10' MINUTE),
  TO_TIMESTAMP(FROM_UNIXTIME(orderTime / 1000, 'yyyy-MM-dd HH:mm:ss'))
)

其中一端是 TIMESTAMP_LTZ、另一端是 TIMESTAMP,需要统一类型。原文给出的做法是把后一个时间值转为 TIMESTAMP_LTZ:


TIMESTAMPDIFF(
  MINUTE,
  (current_timestamp - INTERVAL '10' MINUTE),
  TO_TIMESTAMP_LTZ(mainOrderInfo.orderTime, 3)
)

若确实需要将 CURRENT_TIMESTAMP 转为 TIMESTAMP,文章给出:


TO_TIMESTAMP(CAST(CURRENT_TIMESTAMP AS STRING))

第三类是 Upsert 写入所需的主键。查询产生 update/delete 记录并写入 MySQL、TiDB 等 sink 时,建表语句需声明 PRIMARY KEY,否则会出现 please declare primary key for sink table when query contains update/delete record.。团队尝试利用优化器推断查询主键并自动生成 DDL;原文称其覆盖约 95% 的任务,复杂 SQL 推断失败时会提示用户手动声明。自动推断不等于业务主键一定唯一或符合预期,仍需业务方逐项核验。

原文另展示了 SQL metadata 的 DDL 示例,但该段在文章中出现 CREATEF TABLE 且 MAP 类型不完整,不能直接执行。此处不将错误摘录改写成未经原文支持的可运行代码。

迁移验证与重启顺序

在批量重启生产任务前,团队先为不同任务构建测试任务;次日把测试输出与旧版本实时任务和离线任务对比,以检查数据准确性。SQL 转换后的时间逻辑和自动生成的主键还要让用户再次确认。

正式迁移尽量选择流量较低的时段,按照任务优先级和实时任务血缘规划重启顺序。有赞的做法是先重启低优先级、且位于链路下游的任务;待新版本稳定运行一段时间,再重启更高优先级任务。这样可以降低异常对业务的影响,但并不保证不会出现延迟或数据问题。

升级中遇到的问题

旧 checkpoint 恢复失败

少数 Flink 1.13 作业无法从旧 checkpoint 恢复。作者通过社区邮件和源码排查到,Flink 1.11 后 BaseRowSerializer 更名为 RowDataSerializer,旧状态文件引用的类在新版本中不存在;文章称当时即使使用 State Processor API 也无法处理该缺失类。问题并非每个任务都会触发。迁移前应验证状态恢复、准备数据重放和回退方案;选择凌晨只能减少部分业务影响,不代表恢复一定成功或数据一定不会丢失。

MySQL 维表的无符号 BIGINT 不兼容

Flink 1.13 的 Table connector 统一使用 RowData 后,部分 MySQL 维表关联在字段类型转换时失败:MySQL 的 BIGINT UNSIGNED 不能与 Flink BIGINT 直接相互转换。原文援引 FLINK-18580 给出的建议:在 Flink 维表定义中将该字段声明为 DECIMAL(20,0)。

多条 INSERT 应使用 StatementSet

有赞在 Flink 1.10 中通过 executeSql() 执行同一任务内的多条 INSERT;升级后照旧调用,在 Flink 1.13 中每条 INSERT 会立即提交一个作业并返回相应的 TableResult。团队观察到只有第一条 INSERT 有结果,集群中却提交了多个相同作业。处理方式是将 INSERT 收集进 StatementSet 后统一执行,其他 SQL 仍可用 executeSql():


StatementSet statementSet = streamTableEnv.createStatementSet();
sqls.forEach(sql -> {
  if (isInsertSql(sql)) {
    statementSet.addInsertSql(sql);
  } else {
    streamTableEnv.executeSql(sql);
  }
});
statementSet.execute();

结果与适用边界

文章发表时,有赞报告已将实时计算引擎从 Flink 1.10 升级到 1.13,并完成所有 Flink SQL 任务迁移,稳定运行近三个月;当时 Flink SQL 任务接近实时任务总量的 60%。团队计划继续推进 Kubernetes 化,但文末明确说 Native K8s Application 模式还处于测试阶段。因此,本文不能当作已经完成全量 K8s 迁移的教程。

整体经验是:升级工作远不止替换依赖版本。连接器内部数据结构、UDF、SQL 配置名、时间类型、Upsert 主键、checkpoint 兼容以及多 INSERT 提交行为都需要分别排查,并通过新旧输出对比、状态恢复演练、分批重启和重放预案来控制迁移风险。以上结论限定于有赞当时的 Flink 1.10 至 1.13.2 环境。

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

请登录后发表评论

    暂无评论内容