《Flink 实战与性能优化》1.2 节精读:彻底了解大数据实时计算框架 Flink 的架构、部署与核心特性 示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载本文是 flink-learning 仓库《Flink 实战与性能优化》书中 1.2 节 的系统化展开聚焦于彻底了解 Flink这一主题从数据集类型与运算模型出发逐层拆解 Flink 的整体架构、四种主流部署方式、分布式运行流程Job Client / JobManager / TaskManager、四层 API 抽象、程序与数据流结构、丰富的 Connector 生态以及时间语义、窗口机制、并行执行、状态与容错、内存管理、扩展库等核心特性。读完本篇你将建立 Flink 从是什么到怎么部署、怎么运行、怎么写出一个完整作业的完整认知框架并能结合仓库中对应的源码示例Source / Sink / Window / Checkpoint / State亲手验证每个特性。1 从哪里开始数据集类型与数据运算模型在正式介绍 Flink 之前需要先建立两个基础概念因为整个流式计算框架的设计都围绕它们展开。1.1 数据集类型有界与无界数据集类型分为两大类无界数据集Unbounded无穷的、持续集成的数据集合数据源永不终止有界数据集Bounded有限、不会改变的数据集合处理完即结束。常见的无界数据集包括用户与客户端的实时交互数据应用实时产生的日志金融市场的实时交易记录……1.2 数据运算模型流式与批处理对应的数据运算模型也有两种流式处理只要数据一直在产生计算就持续地进行结果随数据到达不断更新批处理在预先定义的时间范围内运行计算计算完成时释放计算机资源。理解这两组概念是理解流批统一的前提——Flink 的核心设计正是要在这两种模型上收敛。2 Flink 简介真正的流批统一计算引擎Flink 是一个针对流数据和批数据的分布式处理引擎代码主要由 Java 实现部分代码为 Scala。它可以处理有界的批量数据集也可以处理无界的实时数据集。对 Flink 而言其所要处理的主要场景就是流数据批数据只是流数据的一个极限特例因此 Flink 也是一款真正的流批统一的计算引擎。围绕流批统一这一目标Flink 提供了State状态、Checkpoint检查点、Time时间语义、Window窗口等核心机制它们共同构成了 Flink 的基石本文后续将逐项介绍。3 大数据计算引擎的四代演进在网上有人将大数据计算引擎的发展分为四个阶段理解这一演进脉络有助于体会 Flink 出现的背景与定位第一代Hadoop 承载的 MapReduce第二代支持 DAG有向无环图框架的计算引擎 Tez 和 Oozie主要还是批处理任务第三代支持 Job 内部 DAG有向无环图以 Spark 为代表第四代大数据统一计算引擎包括流处理、批处理、AI、Machine Learning、图计算等以 Flink 为代表。或许有人不同意这样的分类其实分类本身并不重要重要的是体会各框架之间的差异与各自更适用的场景。没有任何一个框架可以完美支持所有场景也就不可能存在任何一个框架能完全取代另一个。4 Flink 整体架构自下而上的四层结构Flink 的整体架构从下至上可分为四层部署层Flink 支持本地运行直接在 IDE 中运行程序、在独立集群Standalone 模式或由 YARN、Mesos、K8s 管理的集群上运行也能部署在云上运行层Flink 的核心是分布式流式数据引擎数据以一次一个事件的形式被处理API 层DataStream、DataSet、Table API SQL扩展库层Flink 还包括用于 CEP复杂事件处理、机器学习、图形处理等场景的扩展库。这一分层设计决定了 Flink 既能向上层用户提供从底层原语到高层声明式 SQL 的多种开发选择也能向下层对接多种资源管理系统。5 Flink 的多种方式部署作为一个计算引擎要做得足够完善除了自身特性齐备还必须支持各种生态圈部署方式就是其中之一。Flink 支持以 Local、Standalone、YARN、Kubernetes、Mesos 等形式部署其中最常见的是前四种。5.1 Local 模式直接在 IDE 中运行 Flink Job 时会在本地启动一个 mini Flink 集群适用于开发调试场景。仓库中的大多数示例都采用这种方式例如 并行度示例 通过StreamExecutionEnvironment.getExecutionEnvironment()创建运行环境后直接execute。5.2 Standalone 模式在 Flink 安装目录下执行bin/start-cluster.sh脚本即可启动一个 Standalone 模式的集群。这是 Flink 自带的独立集群模式不依赖外部资源管理系统。5.3 YARN 模式YARN 是 Hadoop 集群的资源管理系统可以在集群上运行各种分布式应用程序。Flink 可与其他应用并行于 YARN 中即 Flink on YARN由 YARN 负责资源申请与回收Flink 作为其上的分布式应用运行。这是生产环境中非常主流的部署方式尤其适合已经建设了 Hadoop 生态的公司。5.4 Kubernetes 模式Kubernetes 是 Google 开源的容器集群管理系统在 Docker 技术基础上为容器化的应用提供部署运行、资源调度、服务发现和动态伸缩等一系列完整功能提高了大规模容器集群管理的便捷性。Flink 也支持部署在 Kubernetes 上。仓库中专门维护了 flink-learning-k8s 模块其中既包含 K8s 部署所需的 Dockerfile 与镜像构建脚本如Dockerfile-flink-1.12.0-jar、Dockerfile-flink-1.12.0-sql、build_flink_docker_images.sh、docker-entrypoint.sh也包含 Flink K8s 相关的 60 个 Java 源码文件以及 Flink K8s 部署博客涵盖 HA 配置、Pod 环境变量、Request 与 Limit 设置、Pod 异常排查等实战话题可对照参考。除以上四种外Flink 还支持 AWS、MapR、Aliyun OSS 等环境。6 Flink 分布式运行流程一次作业提交的背后Flink 作业提交的架构流程由四个关键角色构成理解它们即可掌握 Flink 的分布式执行模型。6.1 Program Code即我们编写的 Flink 应用程序代码定义了数据源、转换逻辑与输出目标。6.2 Job ClientJob Client 不是 Flink 程序执行的内部组成部分但它是任务执行的起点。它负责接收用户的程序代码创建数据流将数据流提交给 JobManager 以便进一步执行执行完成后Job Client 将结果返回给用户。6.3 JobManagerJobManager 是主进程作业管理器协调和管理程序的执行。其主要职责包括安排任务Scheduling、管理 Checkpoint、故障恢复Failover等。集群中至少要有一个 mastermaster 负责调度 task、协调 Checkpoints 和容灾开启高可用HA设置时可以有多个 master但必须保证一个是 leader其余是 standby。JobManager 内部包含Actor system、Scheduler、Check pointing三个重要组件。6.4 TaskManagerTaskManager 从 JobManager 处接收需要部署的 Task是在 JVM 中的一个或多个线程中执行任务的工作节点。任务执行的并行性由每个 TaskManager 上可用的任务槽Slot 个数决定每个任务代表分配给任务槽的一组资源。关于 Slot 有几个关键细节例如 TaskManager 有四个插槽那么它会为每个插槽分配约 25% 的内存每个任务槽中可以运行一个或多个线程同一插槽中的线程共享同一个 JVM同一 JVM 中的任务共享 TCP 连接和心跳消息TaskManager 的一个 Slot 代表一个可用线程该线程拥有固定的内存注意 Slot 只对内存做隔离没有对 CPU 做隔离默认情况下Flink 允许子任务共享 Slot——即使是不同 task 的 subtask只要它们来自同一个 job。这种共享可以带来更好的资源利用率。7 Flink 的 API 抽象层级Flink 提供了不同抽象级别的 API 来开发流式或批处理应用由低到高分为四层最底层有状态流Stateful Stream Processing。通过 Process Function 嵌入到 DataStream API 中允许用户自由处理来自一个或多个流数据的事件并使用具有一致性、容错能力的状态。除此之外用户还可以注册事件时间和处理时间的回调从而实现复杂的计算逻辑DataStream / DataSet API这是 Flink 的核心 API。DataSet 处理有界数据集DataStream 处理有界或无界数据流。用户可以通过 map / flatMap / window / keyBy / sum / max / min / avg / join 等各种方法对数据进行转换或计算Table API以表为中心的声明式 DSL表中的数据可能会动态变化在表达流数据时。Table API 提供了 select、project、join、group-by、aggregate 等操作代码量更少、使用更简洁并且可以在表与 DataStream / DataSet 之间无缝切换允许程序将 Table API 与 DataStream、DataSet 混合使用最高层SQL。语法与表达能力上与 Table API 类似但是以 SQL 查询表达式的形式表现程序。SQL 抽象与 Table API 交互密切SQL 查询可以直接在 Table API 定义的表上执行。为什么 Flink 要通过 SQL 构建统一的大数据流批处理引擎因为在企业中通常存在每天定时生成报表的批处理需求每晚定时跑一遍昨天的数据生成结果报表同时也存在实时性要求很高的流处理需求如实时监控、实时告警。于是整个公司的技术选型越来越多开发人员要学习两套不同的技术框架运维人员也要对两种不同的框架进行环境搭建、作业部署和稳定性维护。当系统越来越复杂、作业越来越多这对开发和运维都是巨大的负担。Flink 通过 SQL API 解决批流统一的痛点后开发和运维只需关注一个计算框架从而减少企业的用人成本与后期开发运维成本。仓库中的 flink-learning-sql 模块含 blink、sql-client、sql-common 三个子模块提供了 Table API SQL 的完整练习代码。8 Flink 程序与数据流结构Source → Transformation → Sink一个完整的 Flink 应用程序由三大部分组成仓库中的示例项目 flink-learning-examples 和 flink-learning-basic 为每一部分都提供了可直接运行的代码。8.1 Source数据输入Flink 在流处理和批处理上的 Source 大致有 4 类基于本地集合的 source基于文件的 source基于网络套接字socket的 source自定义的 source。自定义 Source 中常见的有 Apache Kafka、Amazon Kinesis Streams、RabbitMQ、Twitter Streaming API、Apache NiFi 等当然你也可以实现自己的 Source。仓库中的 Kafka Source 示例 展示了最典型的用法通过Properties配置bootstrap.servers、zookeeper.connect、group.id、key/value.deserializer、auto.offset.reset等参数然后用env.addSource(new FlinkKafkaConsumer(metric, new SimpleStringSchema(), props))从 Kafka topic 中读取数据并print到控制台。该模块下还包含 SourceFromMySQL 等自定义 Source 示例。8.2 Transformation数据转换数据转换的操作种类非常多包括 Map / FlatMap / Filter / KeyBy / Reduce / Fold / Aggregations / Window / WindowAll / Union / Window Join / Split / Select / Project 等可以将数据转换计算成任意想要的结果。仓库 examples/streaming 下按主题划分了丰富的转换示例broadcast广播流、sideoutput旁路输出、joinWindow Join、processFunctionProcess Function 底层 API、async异步 I/O等可直接对照学习。8.3 Sink数据输出Sink 是 Flink 将转换计算后的数据发送的地点通常需要把结果存储下来。Flink 常见的 Sink 大致有写入文件、打印出来print()、写入 socket、自定义 Sink。自定义 Sink 常见的有 Apache Kafka、RabbitMQ、MySQL、ElasticSearch、Apache Cassandra、Hadoop FileSystem 等同样也可以自己定义 Sink。仓库中的 MySQL Sink 示例 是一个非常典型、可直接照抄的自定义 Sink 实现其核心套路是继承RichSinkFunctionTopen()方法中建立数据库连接连接只需要建立一次不必每次 invoke 都重新建立并预编译PreparedStatement每条数据到达时调用invoke()组装参数后执行插入close()方法中关闭连接、释放资源。配套的 数据 Sink 入口类 则展示了完整的作业组装从 Kafka 读取 JSON 字符串经map反序列化为Student对象再addSink(new SinkToMySQL())写入 MySQL。9 丰富的 Connector数据源与存储的生态支撑通过仓库源码可以发现Flink 支持不同版本的 Kafka、不同版本的 ElasticSearch、Cassandra、HBase、Hive、HDFS、RabbitMQ 等 Connector除了支持流作业目前还支持 SQL 作业。由于 Flink 在大数据领域的定位是实时计算、本身不做存储虽然有 State 存储状态数据但这里的存储类似于 MySQL、ElasticSearch 等外部存储所以计算时需要考虑的是数据源来自哪里、计算结果存储到哪里。庆幸的是 Flink 已经支持大部分常用组件。仓库 flink-learning-connectors 目录下按连接器逐个子模块维护了完整示例与文档列举一一对应不同版本的 Kafkaflink-learning-connectors-kafka不同版本的 ElasticSearches5 / es6 / es7 / universalRedisflink-learning-connectors-redisMySQL / JDBCflink-learning-connectors-mysql、flink-learning-connectors-jdbcCassandraflink-learning-connectors-cassandraRabbitMQflink-learning-connectors-rabbitmqHBase1.4 / 2.2flink-learning-connectors-hbaseHDFSflink-learning-connectors-hdfs此外还有 ClickHouse、Flume、Hive、InfluxDB、Kudu、Netty、NiFi、Pulsar、RocketMQ、GCP PubSub、ActiveMQ 等连接器模块除了这些自带的 Connector还可以通过 Flink 提供的接口自定义 Source 和 Sink对应书中的 3.8 节仓库示例可参考上文 8.1 / 8.3 中的实现。10 时间语义Event Time / Ingestion Time / Processing TimeFlink 支持多种 Time 语义包括Event Time事件时间、Ingestion Time摄入时间、Processing Time处理时间这是决定计算结果正确性的关键概念Event Time事件发生的时间由数据本身携带能正确处理乱序与迟到数据需要配合 Watermark水位线使用Ingestion Time事件进入 Flink 系统的时间Processing Time事件被算子处理时的机器系统时间最简单但结果不确定。仓库 examples/streaming/watermark 模块下有 4 个 Main 类与 WordPeriodicWatermark.java、WordPunctuatedWatermark.java分别演示了周期型与定点型 Watermark 的生成方式可对照本书后续 3.1 节深入学习。11 灵活的窗口机制Time / Count / Session / 自定义Flink 支持多种 Window包括Time Window时间窗口、Count Window计数窗口、Session Window会话窗口还支持自定义 Window通过自定义 Trigger、Evictor 等。本书 3.2 节会详细讲解 Window 概念。仓库 flink-learning-window 模块是窗口学习的配套练习其 Main.java 通过 socket 文本流依次演示了滚动时间窗口timeWindow(Time.seconds(30))、滑动时间窗口timeWindow(Time.seconds(60), Time.seconds(30))、计数窗口countWindow(3)、滑动计数countWindow(4, 3)以及会话窗口ProcessingTimeSessionWindows的写法CustomTriggerMain.java 与 CustomTrigger.java 则演示了自定义窗口触发器的实现对应自定义 Window的能力。12 并行执行任务机制Flink 的程序内在是并行和分布式的数据流可以被分区成stream partitions流分区operators 被划分为operator subtasks算子子任务这些 subtasks 在不同的机器或容器中分不同的线程独立运行。operator subtasks 的数量在具体的 operator 上就是其并行计算数Parallelism程序不同的 operator 阶段可能有不同的并行数。例如典型的例子source operator 的并行数为 2而最后的 sink operator 为 1。仓库 并行度示例 直接演示了并行度的设置方式env.setParallelism(1)设置全局并行度map(...).setParallelism(3)则只为某个算子单独设置并行度——这正好印证了程序不同 operator 阶段可能有不同的并行数这一特性。13 状态存储与容错Flink 是一款有状态的流处理框架它提供了丰富的状态访问接口这是它区别于早期无状态流处理框架的核心能力。13.1 状态的分类与数据结构按照数据的划分方式状态可以分为Keyed State键控状态和Operator State算子状态。其中 Keyed State 提供了多种数据结构ValueState单值状态MapStateMap 状态ListState列表状态ReducingState归约状态AggregatingState聚合状态13.2 状态存储方式State Backend状态存储支持多种方式MemoryStateBackend存储在内存中JobManager 的堆内存FsStateBackend存储在文件系统中Checkpoint 持久化到 HDFS 等文件系统RocksDBStateBackend存储在 RocksDB 中本地磁盘支持超大规模状态。13.3 Checkpoint 与 SavepointFlink 通过Checkpoint来提高程序的可靠性开启 Checkpoint 之后Flink 会按照一定的时间间隔对程序的运行状态进行备份当发生故障时Flink 会将所有任务的状态恢复至最后一次 Checkpoint 中的状态并从那里重新开始执行。另外 Flink 还支持通过Savepoint从已停止作业的运行状态进行恢复这种方式需要通过命令手动触发。仓库中的 PvStatExactlyOnce.java 是一个教科书式的状态 容错示例它设置每 1 分钟一次 Checkpoint、语义为 EXACTLY_ONCE并通过enableExternalizedCheckpoints(RETAIN_ON_CANCELLATION)在取消作业时保留外部化 Checkpoint计算逻辑上按 appId 做 KeyBy在RichMapFunction中维护ValueStateLong类型的 pvState每条数据到来时对状态值加 1 后更新——这就是状态支撑精确一次语义的直观体现。同目录下还有 PvStatLocalKeyByExactlyOnce.java 与 PvStatExactlyOnceKafkaUtil.java 可继续深入。此外flink-learning-basic/flink-learning-state 模块提供了 UnionListStateOperator State、Queryable State、状态序列化MetadataSerializer等进阶示例。14 自己的内存管理机制Flink 并不是直接把对象存放在堆内存上而是将对象序列化为固定数量的预先分配的内存段。它采用类似 DBMS 的排序和连接算法可以直接操作二进制数据从而将序列化和反序列化开销降到最低如果需要处理的数据容量超过内存Flink 的运算符会将部分数据存储到磁盘。这种主动内存管理和二进制数据操作带来几个好处保证内存可控可以防止 OutOfMemoryError减少垃圾收集GC压力节省数据的存储空间高效的二进制操作。Flink 是如何分配内存、如何将对象序列化和反序列化、如何操作二进制数据的可以参考仓库 pics/Flink-code.png 所对应的源码解读系列内容如《Flink 是如何管理好内存的》其中讲解了 Flink 的内存管理机制。15 多种扩展库Flink 扩展库中包含了机器学习、Gelly 图形处理、CEP复杂事件处理、State Processing API状态处理 API等这些扩展库在一些特殊场景下会比较适用。仓库 flink-learning-libraries 模块维护了对应的配套练习其中 flink-learning-libraries-cep 提供了 6 个 CEP 复杂事件处理示例flink-learning-libraries-state-processor-api 提供了基于 State Processor API 读写作业状态的示例。这部分内容可对照本书第六章深入学习。16 小结与反思本节在介绍 Flink 之前先讲解了数据集类型有界/无界与数据运算模型流式/批处理随后系统介绍了 Flink 的多个特性整体架构的分层设计、Local/Standalone/YARN/Kubernetes 等部署方式、Program Code → Job Client → JobManager → TaskManager 的分布式运行流程、四层 API 抽象、Source/Transformation/Sink 程序结构、丰富的 Connector 生态以及 Time 语义、Window 机制、并行执行、状态与容错Checkpoint / Savepoint、内存管理、扩展库等。读完本节你不妨带着两个问题回到实践中其一Flink 的哪些特性解决了你在真实业务中遇到的痛点实时性、状态一致性、容灾恢复其二把这些特性与其他计算框架逐一对比你会更清楚什么场景该选择 Flink、什么场景它并不合适。带着这些问题进入本书后续章节Time、Window、状态、自定义 Source/Sink 的深入讲解并在仓库 flink-learning-examples 的示例代码中逐一验证是掌握 Flink 最高效的路径。赞分享示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐flink-learning 开源项目配套专栏《大数据实时计算引擎 Flink 实战与性能优化》内容架构与学习路线全解析flink learning 开源项目配套专栏《大数据实时计算引擎 Flink 实战与性能优化》内容架构与学习路线全解析 《大数据实时计算引擎 Flink 实战示例工程大数据school-of-sre大数据处理框架Flink实时计算与SRE监控school of sre大数据处理框架Flink实时计算与SRE监控 实时计算技术演进与挑战 随着数据生成速度的指数级增长如 UPI交易案例 https:教程MediaCrawler 上手指南快速跑通小红书、抖音、微博等 7 大平台的公开数据爬虫MediaCrawler 上手指南快速跑通小红书、抖音、微博等 7 大平台的公开数据爬虫 MediaCrawler 是一个多平台自媒体数据采集工具支持小红书网页爬虫数据工程上一篇YuE2-3B生产部署终极指南vLLM高并发服务、显存预算与CC BY-NC 4.0商用合规全说明下一篇PyTroll Satpy 开源项目教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考