惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Latest news
Latest news
T
Troy Hunt's Blog
V
Vulnerabilities – Threatpost
L
LINUX DO - 热门话题
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Simon Willison's Weblog
Simon Willison's Weblog
V
V2EX
博客园 - 司徒正美
B
Blog RSS Feed
AWS News Blog
AWS News Blog
MyScale Blog
MyScale Blog
Scott Helme
Scott Helme
Cisco Talos Blog
Cisco Talos Blog
Last Week in AI
Last Week in AI
NISL@THU
NISL@THU
博客园 - Franky
P
Proofpoint News Feed
博客园_首页
C
CERT Recently Published Vulnerability Notes
雷峰网
雷峰网
S
Schneier on Security
P
Proofpoint News Feed
Hugging Face - Blog
Hugging Face - Blog
G
GRAHAM CLULEY
博客园 - 三生石上(FineUI控件)
月光博客
月光博客
WordPress大学
WordPress大学
The Hacker News
The Hacker News
T
Threatpost
阮一峰的网络日志
阮一峰的网络日志
A
Arctic Wolf
Microsoft Azure Blog
Microsoft Azure Blog
T
The Exploit Database - CXSecurity.com
Engineering at Meta
Engineering at Meta
罗磊的独立博客
T
The Blog of Author Tim Ferriss
D
Darknet – Hacking Tools, Hacker News & Cyber Security
I
Intezer
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
K
Kaspersky official blog
SecWiki News
SecWiki News
云风的 BLOG
云风的 BLOG
美团技术团队
C
Cybersecurity and Infrastructure Security Agency CISA
博客园 - 【当耐特】
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
Security Latest
Security Latest
C
Cyber Attacks, Cyber Crime and Cyber Security
B
Blog
S
Security Affairs

Jark's Blog

Flink 1.16:Hive SQL 如何平迁到 Flink SQL Flink CDC 如何简化实时数据入湖入仓 基于 Flink SQL 构建流批一体的 ETL 数据集成 Nexmark: 如何设计一个流计算基准测试? Demo:基于 Flink SQL 构建流式应用 Flink 1.9 实战:使用 SQL 读取 Kafka 并写入 MySQL Flink SQL 编程实践 如何从小白成长为 Apache Committer? 聊聊Blink开源和Flink社区近况 Flink 小贴士 (7): 4个步骤,让 Flink 应用达到生产状态 Flink 小贴士 (6): 使用 Broadcast State 的 4 个注意事项 Flink 小贴士 (5): Savepoint 和 Checkpoint 的 3 个不同点 Flink 小贴士 (4): 如何选择状态后端 Flink 小贴士 (3): 轻松理解 Watermark 一文了解 Apache Flink 核心技术 Flink 零基础实战教程:如何计算实时热门商品 5分钟从零构建第一个 Flink 应用 Flink小贴士 (1):确定Flink作业所需资源大小时要考虑的6件事 Flink在美团的实践与应用 我在阿里的这两年 Flink 原理与实现:Aysnc I/O Flink 原理与实现:Table & SQL API Flink 原理与实现:Session Window Flink 原理与实现:Window 机制 Flink 原理与实现:数据流上的类型和操作 Flink 原理与实现:如何生成 JobGraph Flink 原理与实现:理解 Flink 中的计算资源 Flink 原理与实现:如何生成 StreamGraph Flink 原理与实现:架构和拓扑概览 Flink 原理与实现:内存管理 Flink 原理与实现:如何处理反压问题 Flink官方文档翻译:安装部署(集群模式) Flink官方文档翻译:安装部署(本地模式) 迟到的2015年终总结 高效Macbook开发之道(工具篇) Clojure学习笔记(三):并发与引用 Clojure学习笔记(二):语法 Clojure学习笔记(一):数据结构 Thrift 实践 Thrift 入门 Metrics 是个什么鬼 之入门教程 使用Redis和SQLAlchemy对Scrapy Item去重并存储 使用Scrapy定制可动态配置的爬虫 编程方式下运行 Scrapy spider 读《程序员必读的职业规划书》 Spark 下操作 HBase(1.0.0 新 API) HBase 集群安装部署 Spark On YARN 集群安装部署 Git 常用技能
Flink 小贴士 (2):Flink 如何管理 Kafka 消费位点
WuChong · 2018-11-04 · via Jark's Blog

发表于   |   分类于 Flink   |  

原文:https://data-artisans.com/blog/how-apache-flink-manages-kafka-consumer-offsets
作者:Fabian Hueske, Markos Sfikas
译者:云邪(Jark)

在本周的《Flink Friday Tip》中,我们将结合例子逐步讲解 Apache Flink 是如何与 Apache Kafka 协同工作并确保来自 Kafka topic 的消息以 exactly-once 的语义被处理。

检查点(Checkpoint)是使 Apache Flink 能从故障恢复的一种内部机制。检查点是 Flink 应用状态的一个一致性副本,包括了输入的读取位点。在发生故障时,Flink 通过从检查点加载应用程序状态来恢复,并从恢复的读取位点继续处理,就好像什么事情都没发生一样。你可以把检查点想象成电脑游戏的存档一样。如果你在游戏中发生了什么事情,你可以随时读档重来一次。

检查点使得 Apache Flink 具有容错能力,并确保了即时发生故障也能保证流应用程序的语义。检查点是以固定的间隔来触发的,该间隔可以在应用中配置。

Apache Flink 中实现的 Kafka 消费者是一个有状态的算子(operator),它集成了 Flink 的检查点机制,它的状态是所有 Kafka 分区的读取偏移量。当一个检查点被触发时,每一个分区的偏移量都被存到了这个检查点中。Flink 的检查点机制保证了所有 operator task 的存储状态都是一致的。这里的“一致的”是什么意思呢?意思是它们存储的状态都是基于相同的输入数据。当所有的 operator task 成功存储了它们的状态,一个检查点才算完成。因此,当从潜在的系统故障中恢复时,系统提供了 excatly-once 的状态更新语义。

下面我们将一步步地介绍 Apache Flink 中的 Kafka 消费位点是如何做检查点的。在本文的例子中,数据被存在了 Flink 的 JobMaster 中。值得注意的是,在 POC 或生产用例下,这些数据最好是能存到一个外部文件系统(如HDFS或S3)中。

第一步:

如下所示,一个 Kafka topic,有两个partition,每个partition都含有 “A”, “B”, “C”, ”D”, “E” 5条消息。我们将两个partition的偏移量(offset)都设置为0.

第二步:

Kafka comsumer(消费者)开始从 partition 0 读取消息。消息“A”正在被处理,第一个 consumer 的 offset 变成了1。

第三步:

消息“A”到达了 Flink Map Task。两个 consumer 都开始读取他们下一条消息(partition 0 读取“B”,partition 1 读取“A”)。各自将 offset 更新成 2 和 1 。同时,Flink 的 JobMaster 开始在 source 触发了一个检查点。

第四步:

接下来,由于 source 触发了检查点,Kafka consumer 创建了它们状态的第一个快照(”offset = 2, 1”),并将快照存到了 Flink 的 JobMaster 中。Source 在消息“B”和“A”从partition 0 和 1 发出后,发了一个 checkpoint barrier。Checkopint barrier 用于各个 operator task 之间对齐检查点,保证了整个检查点的一致性。消息“A”到达了 Flink Map Task,而上面的 consumer 继续读取下一条消息(消息“C”)。

第五步:

Flink Map Task 收齐了同一版本的全部 checkpoint barrier 后,那么就会将它自己的状态也存储到 JobMaster。同时,consumer 会继续从 Kafka 读取消息。

第六步:

Flink Map Task 完成了它自己状态的快照流程后,会向 Flink JobMaster 汇报它已经完成了这个 checkpoint。当所有的 task 都报告完成了它们的状态 checkpoint 后,JobMaster 就会将这个 checkpoint 标记为成功。从此刻开始,这个 checkpoint 就可以用于故障恢复了。值得一提的是,Flink 并不依赖 Kafka offset 从系统故障中恢复。

故障恢复

在发生故障时(比如,某个 worker 挂了),所有的 operator task 会被重启,而他们的状态会被重置到最近一次成功的 checkpoint。Kafka source 分别从 offset 2 和 1 重新开始读取消息(因为这是完成的 checkpoint 中存的 offset)。当作业重启后,我们可以期待正常的系统操作,就好像之前没有发生故障一样。如下图所示:

如果想了解更多有关如何最佳地使用 Apache Flink 与 Apache Kafka,以及一些常见问题,可以访问我们这篇文章 Kafka + Flink: A Practical, How-To Guide


hoxis wechat