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

推荐订阅源

U
Unit 42
GbyAI
GbyAI
人人都是产品经理
人人都是产品经理
T
Tor Project blog
Google DeepMind News
Google DeepMind News
The Register - Security
The Register - Security
爱范儿
爱范儿
雷峰网
雷峰网
MongoDB | Blog
MongoDB | Blog
Vercel News
Vercel News
美团技术团队
博客园 - 三生石上(FineUI控件)
D
DataBreaches.Net
L
LangChain Blog
IT之家
IT之家
博客园_首页
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
Last Week in AI
Last Week in AI
博客园 - 【当耐特】
T
Tailwind CSS Blog
M
MIT News - Artificial intelligence
P
Proofpoint News Feed
Hacker News: Ask HN
Hacker News: Ask HN
Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
量子位
Project Zero
Project Zero
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
博客园 - 聂微东
T
Tenable Blog
aimingoo的专栏
aimingoo的专栏
P
Proofpoint News Feed
T
Threat Research - Cisco Blogs
D
Darknet – Hacking Tools, Hacker News & Cyber Security
Cyberwarzone
Cyberwarzone
C
CERT Recently Published Vulnerability Notes
Microsoft Azure Blog
Microsoft Azure Blog
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
博客园 - 叶小钗
Know Your Adversary
Know Your Adversary
L
Lohrmann on Cybersecurity
C
Cisco Blogs
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
Schneier on Security
Schneier on Security
I
InfoQ
P
Privacy & Cybersecurity Law Blog
Spread Privacy
Spread Privacy
Martin Fowler
Martin Fowler
腾讯CDC
S
Security @ Cisco Blogs
F
Fortinet All Blogs

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 小贴士 (3): 轻松理解 Watermark 一文了解 Apache Flink 核心技术 Flink 零基础实战教程:如何计算实时热门商品 5分钟从零构建第一个 Flink 应用 Flink 小贴士 (2):Flink 如何管理 Kafka 消费位点 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 小贴士 (4): 如何选择状态后端
WuChong · 2018-11-21 · via Jark's Blog

发表于   |   分类于 Flink   |  

原文:https://data-artisans.com/blog/stateful-stream-processing-apache-flink-state-backends
作者:Seth Wiesman, Markos Sfikas
译者:云邪(Jark)

本文我们将深入探讨有状态的流处理,更确切地说是 Apache Flink 中不同的状态后端(state backend)。在以下部分,我们将介绍 Apache Flink 的 3 种状态后端,它们的局限性以及根据具体案例需求选择最合适的状态后端。

在有状态的流处理中,当开发人员启用了 Flink 中的 checkpoint 机制,那么状态将会持久化以防止数据的丢失并确保发生故障时能够完全恢复。选择何种状态后端,将决定状态持久化的方式和位置。

Flink 提供了三种可用的状态后端:MemoryStateBackendFsStateBackend,和RocksDBStateBackend

MemoryStateBackend

MemoryStateBackend 是将状态维护在 Java 堆上的一个内部状态后端。键值状态和窗口算子使用哈希表来存储数据(values)和定时器(timers)。当应用程序 checkpoint 时,此后端会在将状态发给 JobManager 之前快照下状态,JobManager 也将状态存储在 Java 堆上。默认情况下,MemoryStateBackend 配置成支持异步快照。异步快照可以避免阻塞数据流的处理,从而避免反压的发生。

使用 MemoryStateBackend 时的注意点:

  • 默认情况下,每一个状态的大小限制为 5 MB。可以通过 MemoryStateBackend 的构造函数增加这个大小。
  • 状态大小受到 akka 帧大小的限制,所以无论怎么调整状态大小配置,都不能大于 akka 的帧大小。也可以通过 akka.framesize 调整 akka 帧大小(通过配置文档了解更多)。
  • 状态的总大小不能超过 JobManager 的内存。

何时使用 MemoryStateBackend:

  • 本地开发或调试时建议使用 MemoryStateBackend,因为这种场景的状态大小的是有限的。
  • MemoryStateBackend 最适合小状态的应用场景。例如 Kafka consumer,或者一次仅一记录的函数 (Map, FlatMap,或 Filter)。

FsStateBackend

FsStateBackend 需要配置的主要是文件系统,如 URL(类型,地址,路径)。举个例子,比如可以是:

  • “hdfs://namenode:40010/flink/checkpoints” 或
  • “s3://flink/checkpoints”

当选择使用 FsStateBackend 时,正在进行的数据会被存在 TaskManager 的内存中。在 checkpoint 时,此后端会将状态快照写入配置的文件系统和目录的文件中,同时会在 JobManager 的内存中(在高可用场景下会存在 Zookeeper 中)存储极少的元数据。

默认情况下,FsStateBackend 配置成提供异步快照,以避免在状态 checkpoint 时阻塞数据流的处理。该特性可以实例化 FsStateBackend 时传入 false 的布尔标志来禁用掉,例如:

new FsStateBackend(path, false);

使用 FsStateBackend 时的注意点:

  • 当前的状态仍然会先存在 TaskManager 中,所以状态的大小不能超过 TaskManager 的内存。

何时使用 FsStateBackend:

  • FsStateBackend 适用于处理大状态,长窗口,或大键值状态的有状态处理任务。
  • FsStateBackend 非常适合用于高可用方案。

RocksDBStateBackend

RocksDBStateBackend 的配置也需要一个文件系统(类型,地址,路径),如下所示:

  • “hdfs://namenode:40010/flink/checkpoints” 或
  • “s3://flink/checkpoints”

RocksDB 是一种嵌入式的本地数据库。RocksDBStateBackend 将处理中的数据使用 RocksDB 存储在本地磁盘上。在 checkpoint 时,整个 RocksDB 数据库会被存储到配置的文件系统中,或者在超大状态作业时可以将增量的数据存储到配置的文件系统中。同时 Flink 会将极少的元数据存储在 JobManager 的内存中,或者在 Zookeeper 中(对于高可用的情况)。RocksDB 默认也是配置成异步快照的模式。

使用 RocksDBStateBackend 时的注意点:

  • RocksDB 支持的单 key 和单 value 的大小最大为每个 2^31 字节。这是因为 RocksDB 的 JNI API 是基于 byte[] 的。
  • 我们需要强调的是,对于使用具有合并操作的状态的应用程序,例如 ListState,随着时间可能会累积到超过 2^31 字节大小,这将会导致在接下来的查询中失败。

何时使用 RocksDBStateBackend:

  • RocksDBStateBackend 最适合用于处理大状态,长窗口,或大键值状态的有状态处理任务。
  • RocksDBStateBackend 非常适合用于高可用方案。
  • RocksDBStateBackend 是目前唯一支持增量 checkpoint 的后端。增量 checkpoint 非常使用于超大状态的场景。

当使用 RocksDB 时,状态大小只受限于磁盘可用空间的大小。这也使得 RocksDBStateBackend 成为管理超大状态的最佳选择。使用 RocksDB 的权衡点在于所有的状态相关的操作都需要序列化(或反序列化)才能跨越 JNI 边界。与上面提到的堆上后端相比,这可能会影响应用程序的吞吐量。

不同状态后端满足不同场景的需求,在开始开发应用程序之前应该仔细考虑和规划后选择。这可确保选择了正确的状态后端以最好地满足应用程序和业务需求。


hoxis wechat