0

0

Flink 1.16 JobManager 重启后消息丢失问题排查与解决

霞舞

霞舞

发布时间:2025-10-08 10:48:39

|

583人浏览过

|

来源于php中文网

原创

flink 1.16 jobmanager 重启后消息丢失问题排查与解决

在 Flink 1.16 中,JobManager 重启后消息丢失是一个比较棘手的问题。以下将从多个角度分析可能的原因,并提供相应的解决方案。

首先,我们引用上面的摘要:本文针对 Flink 1.16 中遇到的 JobManager 重启后消息丢失问题,提供了一系列可能的排查方向和解决方案。文章涵盖了从检查是否陷入死循环、确认 Source 是否支持 Checkpointing 和 Rewindable,到排除 JobManagerCheckpointStorage 导致的 Checkpoint 丢失等多个方面,并提供了高可用配置的建议,旨在帮助读者全面理解并解决此类问题,确保 Flink 作业的稳定性和数据完整性。

1. 检查是否陷入 "Fail -> Restart -> Fail Again" 死循环

最常见的原因是程序遇到了“毒丸”(Poison Pill)数据,即无法被正确处理的记录。Flink 在遇到这种数据时,会不断地尝试重启并重新处理该数据,从而陷入死循环,导致后续消息无法被处理。

解决方案:

  • 数据清洗 在 Source 端进行数据清洗,过滤掉不符合格式或会导致异常的数据。
  • 容错处理: 在算子中添加容错处理逻辑,例如使用 try-catch 捕获异常,并记录错误信息,避免程序崩溃。
  • 侧输出流: 将无法处理的数据发送到侧输出流,进行后续分析和处理。

示例代码:

DataStream input = ...;

DataStream processed = input.map(value -> {
    try {
        // 尝试处理数据
        return process(value);
    } catch (Exception e) {
        // 记录错误信息
        LOG.error("Error processing value: {}", value, e);
        // 将数据发送到侧输出流
        return null; // 或者抛出异常,并使用侧输出流捕获
    }
}).filter(Objects::nonNull); // 过滤掉 null 值

// 获取侧输出流
OutputTag errorTag = new OutputTag("error-tag", Types.STRING);
DataStream errorStream = processed.getSideOutput(errorTag);

2. Source 是否支持 Checkpointing 和 Rewindable

Flink 的容错机制依赖于 Checkpointing 和 Source 的 Rewindable 能力。如果 Source 不支持 Checkpointing,或者无法 Rewind 到上次 Checkpoint 的位置,则会导致数据丢失

解决方案:

  • 选择支持 Checkpointing 的 Source: 尽量选择 Flink 官方提供的或经过验证的、支持 Checkpointing 的 Source Connector。
  • 自定义 Source: 如果必须使用不支持 Checkpointing 的 Source,可以考虑自定义 Source,并实现 Checkpointing 接口。 需要注意的是,自定义 Source 的 Checkpointing 实现较为复杂,需要仔细考虑数据一致性和性能问题。
  • 使用 Kafka 或其他可靠消息队列: 将数据先写入 Kafka 等可靠消息队列,再从 Kafka 读取数据进行处理。Kafka 具有持久化存储和回溯消费的能力,可以保证数据不丢失。

3. Checkpoint 存储位置

Checkpoint 的存储位置也会影响数据恢复。如果使用 JobManagerCheckpointStorage,则 Checkpoint 数据存储在 JobManager 的内存中。当 JobManager 重启时,Checkpoint 数据会丢失,导致数据无法恢复。

Tome
Tome

先进的AI智能PPT制作工具

下载

解决方案:

  • 配置高可用性存储: 使用高可用性的 Checkpoint 存储,例如 HDFS、RocksDB 或 S3。这些存储方式可以将 Checkpoint 数据持久化存储,即使 JobManager 重启,数据也不会丢失。

配置示例 (flink-conf.yaml):

state.checkpoints.dir: hdfs:///flink/checkpoints
state.savepoints.dir: hdfs:///flink/savepoints
state.backend: rocksdb
state.backend.rocksdb.memory.managed: true

4. JobManager 高可用配置

JobManager 的重启不应该导致数据丢失。如果 JobManager 重启导致数据丢失,说明集群的高可用配置可能存在问题。

解决方案:

  • 配置 ZooKeeper 或其他高可用协调服务: 配置 ZooKeeper 或其他高可用协调服务,用于选举 Leader JobManager,并在 Leader JobManager 失败时自动切换到备用 JobManager。
  • 配置高可用性存储: 如前所述,使用高可用性的 Checkpoint 存储,确保 Checkpoint 数据不丢失。

高可用配置示例 (flink-conf.yaml):

high-availability: zookeeper
high-availability.storageDir: hdfs:///flink/ha
high-availability.cluster-id: /flink-cluster
high-availability.zookeeper.quorum: zk-host1:2181,zk-host2:2181,zk-host3:2181

总结

解决 Flink JobManager 重启后消息丢失问题,需要从多个方面进行排查。首先要确定是否陷入死循环,然后检查 Source 是否支持 Checkpointing 和 Rewindable,接着排除 Checkpoint 存储位置的影响,最后配置 JobManager 的高可用性。通过以上步骤,可以有效地解决该问题,保证 Flink 作业的稳定性和数据完整性。

注意事项:

  • 在生产环境中,务必配置高可用性的 Checkpoint 存储和 JobManager,以确保数据安全。
  • 定期检查 Flink 集群的日志,及时发现并解决潜在问题。
  • 监控 Flink 作业的运行状态,及时发现并处理异常情况。
  • 升级到最新的 Flink 版本,可以获得更好的性能和稳定性。

热门AI工具

更多
DeepSeek
DeepSeek

幻方量化公司旗下的开源大模型平台

豆包大模型
豆包大模型

字节跳动自主研发的一系列大型语言模型

通义千问
通义千问

阿里巴巴推出的全能AI助手

腾讯元宝
腾讯元宝

腾讯混元平台推出的AI助手

文心一言
文心一言

文心一言是百度开发的AI聊天机器人,通过对话可以生成各种形式的内容。

讯飞写作
讯飞写作

基于讯飞星火大模型的AI写作工具,可以快速生成新闻稿件、品宣文案、工作总结、心得体会等各种文文稿

即梦AI
即梦AI

一站式AI创作平台,免费AI图片和视频生成。

ChatGPT
ChatGPT

最最强大的AI聊天机器人程序,ChatGPT不单是聊天机器人,还能进行撰写邮件、视频脚本、文案、翻译、代码等任务。

相关专题

更多
kafka消费者组有什么作用
kafka消费者组有什么作用

kafka消费者组的作用:1、负载均衡;2、容错性;3、广播模式;4、灵活性;5、自动故障转移和领导者选举;6、动态扩展性;7、顺序保证;8、数据压缩;9、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

168

2024.01.12

kafka消费组的作用是什么
kafka消费组的作用是什么

kafka消费组的作用:1、负载均衡;2、容错性;3、灵活性;4、高可用性;5、扩展性;6、顺序保证;7、数据压缩;8、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

151

2024.02.23

rabbitmq和kafka有什么区别
rabbitmq和kafka有什么区别

rabbitmq和kafka的区别:1、语言与平台;2、消息传递模型;3、可靠性;4、性能与吞吐量;5、集群与负载均衡;6、消费模型;7、用途与场景;8、社区与生态系统;9、监控与管理;10、其他特性。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

202

2024.02.23

硬盘接口类型介绍
硬盘接口类型介绍

硬盘接口类型有IDE、SATA、SCSI、Fibre Channel、USB、eSATA、mSATA、PCIe等等。详细介绍:1、IDE接口是一种并行接口,主要用于连接硬盘和光驱等设备,它主要有两种类型:ATA和ATAPI,IDE接口已经逐渐被SATA接口;2、SATA接口是一种串行接口,相较于IDE接口,它具有更高的传输速度、更低的功耗和更小的体积;3、SCSI接口等等。

1155

2023.10.19

PHP接口编写教程
PHP接口编写教程

本专题整合了PHP接口编写教程,阅读专题下面的文章了解更多详细内容。

213

2025.10.17

php8.4实现接口限流的教程
php8.4实现接口限流的教程

PHP8.4本身不内置限流功能,需借助Redis(令牌桶)或Swoole(漏桶)实现;文件锁因I/O瓶颈、无跨机共享、秒级精度等缺陷不适用高并发场景。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

1907

2025.12.29

java接口相关教程
java接口相关教程

本专题整合了java接口相关内容,阅读专题下面的文章了解更多详细内容。

22

2026.01.19

dubbo和zookeeper有什么区别
dubbo和zookeeper有什么区别

dubbo和zookeeper的区别:1、功能定位;2、使用场景;3、数据存储与协调;4、集成与关系;5、性能与可靠性;6、扩展性与灵活性;7、社区与生态系统。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

226

2024.02.23

C++ 设计模式与软件架构
C++ 设计模式与软件架构

本专题深入讲解 C++ 中的常见设计模式与架构优化,包括单例模式、工厂模式、观察者模式、策略模式、命令模式等,结合实际案例展示如何在 C++ 项目中应用这些模式提升代码可维护性与扩展性。通过案例分析,帮助开发者掌握 如何运用设计模式构建高质量的软件架构,提升系统的灵活性与可扩展性。

9

2026.01.30

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
10分钟--Midjourney创作自己的漫画
10分钟--Midjourney创作自己的漫画

共1课时 | 0.1万人学习

Midjourney 关键词系列整合
Midjourney 关键词系列整合

共13课时 | 0.9万人学习

AI绘画教程
AI绘画教程

共2课时 | 0.2万人学习

关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号 技术交流群
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn | 湘ICP备2023035733号