ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink入门到实践:核心概念、部署架构与实时计算调优

Flink入门到实践:核心概念、部署架构与实时计算调优 1. 先搞清楚Flink 到底解决什么问题1.1 从一个半夜被叫醒的报表说起很多做数据开发的兄弟应该都有过这样的经历业务方第二天一早要出昨天的数据报表跑批脚本凌晨排队跑完还要检查日志稍微遇到点脏数据整张报表就挂了。后来换上了 Flink把离线批处理的流程改成了实时流式处理数据从业务库变化到报表可见从“明天早上见”变成了“秒级可见”半夜被叫醒的情况少了一大半。Apache Flink 本质上是一个分布式处理引擎专为无界和有界数据流上的有状态计算设计。句子里有两个关键词一个是“流”一个是“状态”。流意味着数据是持续的、不断到达的状态意味着同一个 key 的处理需要跨事件、跨时间记住上下文。比如统计每个用户的实时消费总额你得记住每个用户到目前为止累计花了多少钱这个“记住”就是状态。对于初学者第一章最重要的是建立框架思维而不是急着啃源码。你需要先理解 Flink 在解决什么问题、它跟普通批处理框架有何不同、一个作业的基本长什么样。有了这个心智模型后面的安装、部署、调优、面试题才会有落点。1.2 有界数据和无界数据要理解 Flink必须先分清两种数据形态。有界数据就是有开始有结束的数据集合。离线 Hive 表里的一批日志、某个时间段的订单快照都是典型的例子。批处理框架的核心逻辑是“先把所有数据准备好再统一计算”所以结果天然带有延迟适合对时效性要求不高的场景。无界数据是源源不断产生的数据流。用户点击、传感器上报、交易流水这些数据没有终点你永远等不到“全部数据到齐”的那一天。流处理框架的核心逻辑是“数据来了就处理处理完接着等下一条”所以结果可以做到秒级甚至毫秒级延迟。Flink 比较厉害的地方是它把这两种形态统一了。批处理可以看作是流处理的一种特殊形式把有界数据当成一条有限流用同一套引擎处理。这也是 Flink 官方强调“Unified for both stream and batch”的原因。第一章学完你最好能用自己的话说清楚有界和无界的区别以及 Flink 在两种场景下分别怎么处理。1.3 Flink 的不可替代性体现在哪里市面上的流处理框架不止一个Flink 能站住脚靠的是几个硬指标。第一真正做到了流式处理原生化。很多早期大数据的方案是先落盘再算本质上还是“小批”处理Flink 从底层设计上就是事件驱动的每条数据到达后立即触发计算延迟可以压到极低。第二状态管理做得很扎实。分布式流计算如果没有可靠的状态机制重启一次就丢数据那根本没法上生产。Flink 自带状态存储、状态快照、自动恢复配合 Checkpoint 机制可以把状态做得非常稳。第三精确一次语义Exactly-once不再是宣传口号。Flink 通过分布式快照加上两阶段提交在很多外部系统上真的实现了端到端的不重不丢这在金融交易、风控计费场景里是刚需。第四生态和吞吐都很能打。连接器覆盖 Kafka、MySQL、Elasticsearch、Hive 等常用系统SQL 支持也完善很多团队可以直接用 Flink SQL 替代一部分 ETL 工作大幅降低开发成本。所以如果你正在做实时数仓、实时风控、实时推荐、数据同步Flink 基本是绕不开的核心组件。第一章把它的定位搞清楚后面学起来会顺很多。2. 核心架构与第一章必须建立的心智模型2.1 先看分层别一头扎进源码我第一次接触 Flink 的时候被一堆名词砸晕了JobManager、TaskManager、Slot、Operator、State、Checkpoint。后来发现只要先把分层结构记住所有名词各归其位思路一下就清楚了。从使用者的角度俯视 Flink可以把整个体系分成几个层面。最上层是 API 层包括 DataStream API、DataSet API、Table/SQL API。日常开发主要在这一层写逻辑。再往下是运行时层也就是真正干活的引擎负责把代码编译成执行图、调度任务、分配资源、处理故障。再往下是状态和检查点层这是 Flink 的看家本领状态存储、增量快照、恢复机制都在这一层。最底下是部署层支持本地、Standalone、Kubernetes、YARN、云托管等多种模式。我建议初学者记住一条主线用 API 写业务让运行时帮你在分布式集群上跑起来靠状态和检查点保底最终的结果写到外部系统。后面你在 Web UI 上看到的各种指标、异常、反压基本都是围绕这条主线展开的。简单画一个逻辑关系就是Source读数据比如从 Kafka、Kinesis、文件或 Socket 读。Transformation转换计算比如过滤、聚合、窗口、维度关联。Sink输出结果比如写 MySQL、Elasticsearch、Hive、日志。这一套 Source-Transformation-Sink 的 Pipeline 模型在任何流处理框架里都适用。理解了它你就理解了 Flink 作业的骨架。2.2 一个 Flink 作业的基本骨架Flink 的 DataStream 作业写起来非常像一条流水线。看一个最经典的 WordCount 流式版本import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; public class StreamWordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamString lines env.socketTextStream(localhost, 9999); DataStreamTuple2String, Integer counts lines .flatMap((String line, CollectorTuple2String, Integer out) - { for (String word : line.split(\\s)) { out.collect(Tuple2.of(word, 1)); } }) .returns(Types.TUPLE(Types.STRING, Types.INT)) .keyBy(value - value.f0) .sum(1); counts.print(); env.execute(stream-word-count); } }这段代码干了这么几件事通过socketTextStream从本地 9999 端口读文本行然后把每一行按空格拆成单词每个单词变成(单词, 1)的二元组再按单词分组最后把相同单词的数量累加起来打印到控制台最后用execute把任务提交起来。注意最后一行env.execute非常关键。很多新手一开始忘了写程序直接没有反应因为执行环境没有被真正触发。Flink 是惰性求值前面定义的一系列转换操作都只是构建计算图只有调用了execute()作业才会被真正提交和启动。另一个容易踩的坑是 lambda 表达式需要显式指定返回类型。第一次见.returns(Types.TUPLE(...))可能很奇怪但这是因为泛型在运行时会被类型擦除Flink 无法自动推断出Tuple2String, Integer的真实类型。没有这行运行时会直接报类型相关的异常。2.3 为什么不是自己开一堆线程去算有人可能会想一个词频统计我用多线程加 Map 不也能实现吗为什么要引入 Flink 这么重的框架差距主要藏在分布式环境和故障处理里。单机多线程只能处理单台机器的数据量一旦数据量上去机器内存不够、CPU 打满系统就会崩溃。Flink 解决的是把任务分散到多台机器上并行计算并且让整体表现像一台机器一样一致。更关键的是故障恢复。分布式环境下网络抖动、机器宕机、容器重启都是常态如果每个子任务独自在内存里维护状态一旦某个节点挂了那部分数据就丢了。Flink 通过状态快照把每个算子的状态定期保存到外部存储故障后从最近一次快照恢复做到数据尽量不丢。这个是自研多线程方案几乎不可能在短时间内做出来的能力。所以能用 Flink 解决的问题不要重复造轮子。它虽然学习曲线有点陡但投入产出比非常高。3. 从零部署Linux 和 Docker 两种安装方式3.1 版本选择和 JDK 环境我建议新手直接从 Flink 1.17 或 1.18 版本开始学不选太老的版本。旧版不仅 API 有差异网上教程也多是针对新版本的写法版本对不上会平添不少烦恼。Flink 依赖 JDK一般推荐 JDK 8 或 JDK 11。如果你的机器装的是 JDK 17Flink 1.17 之后也能跑但个别老版本的连接器可能不兼容。所以最稳妥的组合是 Flink 1.17 JDK 8或者 Flink 1.18 JDK 11。官方对 Scala 版本也有区分。Flink 本身用 Java 写成但很多周边代码兼容 Scala 2.12 和 2.13下载包名里会看到bin-scala_2.12这样的后缀。日常用 Java 开发不需要关心 Scala 版本直接用默认的 Scala 2.12 包就行。3.2 Linux 本地安装步骤假设你有一台干净的 Linux 服务器比如 Ubuntu 22.04 或者 CentOS 7可以按下面的步骤操作。先确认 Java 已安装java -version如果没有用包管理器装一个 OpenJDKsudo apt update sudo apt install -y openjdk-11-jdk然后下载 Flink 并解压cd /opt wget https://archive.apache.org/dist/flink/flink-1.17.2/flink-1.17.2-bin-scala_2.12.tgz tar -xzf flink-1.17.2-bin-scala_2.12.tgz cd flink-1.17.2解压完的目录结构里有几个东西要记住bin启动脚本所在目录比如start-cluster.sh、flink。conf配置文件所在目录核心是flink-conf.yaml。lib存放依赖 jar连接器驱动也放这里。log运行日志目录排查问题第一步先看这里。examples官方自带示例 jar。启动集群./bin/start-cluster.sh启动完成后用jps查看进程应该能看到两个 Java 进程一个叫 StandaloneSessionClusterEntrypoint对应 JobManager一个叫 TaskManagerRunner对应 TaskManager。然后用浏览器打开http://localhost:8081就能看到 Flink Web UI。JobManager 是调度中心负责接收作业、分配任务、故障恢复TaskManager 是干活的人真正执行算子逻辑。最后关集群./bin/stop-cluster.sh这套本地模式适合学习和代码调试但不适合生产。生产环境一般会直接部署到 Kubernetes 或使用平台提供的托管版本思路类似只是资源管理方式不同。3.3 用 Docker 快速起一个 Flink 环境如果不想污染本机环境或者你想快速复现别人的示例Docker 是最方便的。Linux 下先确认 Docker 装好再直接拉官方镜像。最简单的方式是直接用 Docker 命令起一个 Flink 集群。先建一个专用网络让 JobManager 和 TaskManager 互相能通信docker network create flink-net docker run -d \ --name jobmanager \ --network flink-net \ -p 8081:8081 \ -e FLINK_PROPERTIESjobmanager.rpc.address: jobmanager \ flink:1.17.2 jobmanager docker run -d \ --name taskmanager \ --network flink-net \ -p 6122:6122 \ -e FLINK_PROPERTIESjobmanager.rpc.address: jobmanager \ flink:1.17.2 taskmanager需要注意FLINK_PROPERTIES里的jobmanager.rpc.address必须指向 JobManager 的容器名否则 TaskManager 找不到 JobManager。如果忘了设置常见现象是 TaskManager 一直在等连接Web UI 里的 TaskManager 列表是空的。如果你经常用 Docker Compose 管理服务我建议把配置写成文件这样团队协作也方便。一个最简的docker-compose.yml可以长这样services: jobmanager: image: flink:1.17.2 container_name: jobmanager ports: - 8081:8081 environment: - FLINK_PROPERTIES jobmanager.rpc.address: jobmanager - | taskmanager.numberOfTaskSlots: 2 command: jobmanager taskmanager: image: flink:1.17.2 container_name: taskmanager depends_on: - jobmanager environment: - FLINK_PROPERTIES jobmanager.rpc.address: jobmanager command: taskmanager启动docker compose up -d浏览器访问http://localhost:8081看到 Web UI 就说明跑起来了。3.4 把 WordCount 作业提交上去环境起来了总得跑一个作业验一验。官方自带了一个示例 jar./bin/flink run \ -m localhost:8081 \ -d \ examples/streaming/WordCount.jar-m指定 JobManager 的地址-d表示后台运行不带-d会一直等待作业结束。看到一个Job has been submitted successfully的提示同时 Web UI 的 Running Jobs 里出现一个作业就说明提交成功了。点击作业进去可以看到每个算子的并行度、输入输出记录数、延迟等指标。第一次跑通这个流程你大概就对 Flink 的“提交-调度-执行”这条链路有了直观感受。如果你想把编写的代码打包提交通常用 Maven 打包生成一个 fat jarmvn clean package ./bin/flink run -m localhost:8081 -d target/flink-demo-1.0.jar这里注意一个细节打包时 Flink 核心依赖的作用域最好设置成provided让 Flink 运行时提供否则 jar 会非常大而且可能引入依赖冲突导致提交后各种 “NoClassDefFoundError”。3.5 部署配置里容易被忽略的几个参数部署第一步虽然简单但配置文件的坑不少。新手最常忽略的是内存参数。Flink 默认的 JobManager 内存是 1600M 左右TaskManager 的内存默认值在不同版本不一样。如果你的服务器内存不大起多个 TaskManager 很容易直接把机器内存打满然后系统开始疯狂 swap作业性能一落千丈。建议在conf/flink-conf.yaml里显式设置jobmanager.memory.process.size: 1024m taskmanager.memory.process.size: 2048m taskmanager.memory.managed.fraction: 0.4 taskmanager.numberOfTaskSlots: 2taskmanager.memory.managed.fraction是给排序、状态后端、窗口缓存使用的堆外内存比例调太大留给用户代码和网络缓冲的内存就少调太小又可能导致 RocksDB 状态后端内存不够。大部分场景先给 0.4 是一个比较安全的起步值。还有一点生产环境千万不要直接改完配置就重启最好加-D参数覆盖或者通过配置中心动态下发避免手滑把集群搞挂。4. 第一章最关键的几个概念时间、窗口与状态4.1 处理时间和事件时间Flink 里有两个时间概念初学者必须分清楚。处理时间Processing Time是指数据到达 Flink 时的机器时间。它简单、实时性高但结果是不可重复的。比如你统计过去一分钟收到的订单如果网络抖动导致一条数据延迟了十秒它的归属时间就变了同一份数据在不同时间跑出来的结果可能不一致。事件时间Event Time是指数据产生时附带的时间戳。比如订单日志里的create_time用户点击日志里的click_time。事件时间更能反映业务真实语义但要应对乱序和延迟问题需要引入水位线机制。日常业务中大部分统计需求都应该优先使用事件时间尤其是做报表、风控、对账这种对准确性要求高的场景。处理时间一般只用于实时监控告警这种不需要精确回放的地方。代码里设置事件时间通常在建 source 之后要指定时间戳提取器和水位线生成器DataStreamOrder orders env .addSource(kafkaSource) .assignTimestampsAndWatermarks( WatermarkStrategy.OrderforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) - event.getCreateTime()) );这段代码的意思是允许事件最多乱序 5 秒每条数据的时间戳从createTime字段提取。Flink 收到事件时间小于当前水位线的数据就会判定为迟到数据默认会被窗口丢弃或进入侧输出流。4.2 水位线到底是个什么东西水位线是 Flink 里理解门槛最高的概念之一很多人会绕晕。我的理解是水位线等于一个“事件时间的边界承诺”。每来一条数据Flink 会基于数据处理逻辑算出当前应该推进到的水位线实际上告诉下游算子等于和小于这个时间戳的事件都已经“基本到齐了”可以放心触发计算。举个例子。订单流的创建时间分别是 10:00:01、10:00:03、10:00:02最后那条是乱序的。如果窗口是 10:00:00 到 10:00:05水位线推进到 10:00:05 之后窗口就会触发计算哪怕还有迟到的 10:00:04 数据没到窗口也不会无限等待。设置水位线越保守窗口触发越晚正确性更高但实时性变差。写法上forBoundedOutOfOrderness(Duration.ofSeconds(5))表示最多容忍 5 秒乱序具体值要根据业务容忍度来调节。线上有个常见做法是结合延迟监控动态调整而不是拍脑袋定一个值。面试的时候能被问到的基本都是水位线是什么、怎么产生、怎么推进、迟到数据怎么处理。能把上面这段话用自己的话讲清楚这一关基本就过了。4.3 窗口怎么开无界流本身没有边界但业务上总得切出一个个范围来统计这就是窗口机制。三种最常用的窗口类型要记牢滚动窗口Tumbling Window固定长度无重叠比如每 5 分钟统计一次。滑动窗口Sliding Window固定长度但可以重叠比如每 5 分钟计算过去 15 分钟的数据。会话窗口Session Window根据不活跃间隔切分比如用户连续 30 分钟没操作就结束一次会话。在 DataStream API 里开窗口一般是这样orders .keyBy(order - order.getUserId()) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new OrderAmountAggregate())在 Flink SQL 里更简单SELECT user_id, SUM(amount) FROM orders GROUP BY user_id, TUMBLE(create_time, INTERVAL 5 MINUTE);新手最容易犯的错误是把TABLE里的WATERMARK FOR create_time AS create_time - INTERVAL 5 SECOND写错字段或者把窗口别名用到分组外。记住一个原则窗口列必须出现在 GROUP BY 中且在 SQL 里窗口函数生成的是一个伪列可以用来投射窗口开始和结束时间。4.4 状态与检查点状态是 Flink 里让流处理“变聪明”的东西。没有状态你只能对当前这一条数据做无脑转换有了状态你才能做累计、去重、关联、窗口聚合。Flink 的状态分为两种托管状态和原始状态。日常开发用托管状态就够了它由 Flink 自动管理存储、恢复和重分布。托管状态又分为 Keyed State 和 Operator State。Keyed State 绑定每一个 key比如每个用户的累计消费金额Operator State 绑定一个算子实例比如记录 source 读到哪了。为了让状态在故障时不丢Flink 会定期做检查点Checkpoint。检查点的原理是基于分布式快照把所有算子的状态和当前处理到的数据位置统一保存一份快照。任务挂掉后从最近一次成功的检查点恢复。生产上需要重点配置几个参数execution.checkpointing.interval: 60s execution.checkpointing.mode: EXACTLY_ONCE state.backend.type: rocksdb检查点间隔太短会导致频繁做快照性能下降太长会导致故障恢复丢的数据多。一般是 30 秒到 5 分钟之间根据业务容忍度调整。状态后端方面如果状态规模大选 RocksDB 更稳如果状态不大且要求性能用 HashMap 状态后端就够了。5. 真实环境里的两件烦心事JDBC 连接器和火焰图5.1 JDBC 连接器常见异常排查Flink 要读写 MySQL、PostgreSQL 这类数据库很多人第一反应是去写 JDBC 连接代码其实连接数据库也有现成连接器。通过 Flink SQL 定义一张 JDBC 维表或结果表就可以直接读写外部数据库不需要自己管理连接池。定义一张 JDBC 表DDL 大概长这样CREATE TABLE user_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), create_time TIMESTAMP(3), WATERMARK FOR create_time AS create_time - INTERVAL 5 SECOND ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/shop, driver com.mysql.cj.jdbc.Driver, table-name t_order, username root, password 123456 );这里面的driver参数一定不能写错。不同版本的 MySQL 驱动类名不一样8.x 是com.mysql.cj.jdbc.Driver5.x 用的是com.mysql.jdbc.Driver。写错之后Flink 会在运行时报 ClassNotFound 或驱动初始化异常。新手最常见的报错我列一个速查表报错现象根本原因解决方式Could not find any factory for identifier jdbc缺 JDBC 连接器 jar把 flink-connector-jdbc 放进 Flink lib 目录ClassNotFoundException: com.mysql.cj.jdbc.Driver缺 MySQL 驱动 jar把 mysql-connector-java 放进 lib 目录Communications link failure网络不通或账号权限不对检查 IP、端口、用户名、白名单Connection is not available, request timed out数据库连接池满或慢查询调大连接池、优化 SQL、减少并发写入sink buffer flush timeout写入频率低导致缓冲堆积配置 sink.buffer-flush.max-rows 和 interval实际操作中Flink SQL 客户端启动后默认不会加载用户自定义目录下的 jar。如果使用sql-client.sh建议加-j参数指定依赖 jar./bin/sql-client.sh \ -j lib/flink-connector-jdbc-3.1.2-1.17.jar \ -j lib/mysql-connector-j-8.0.33.jar或者直接把 jar 放到lib目录下再启动。我之前在这个坑里耗了大半天报错一直提示找不到连接器其实只是 jar 没加载进去。另外JDBC 维表关联时如果底表数据量大查询频繁一定要开启 lookup 缓存lookup.cache.max-rows 1000, lookup.cache.ttl 60s不开启缓存的话每条流数据都会打一次数据库数据库很容易被打垮。不过缓存也会带来脏读数据更新频率高的表要谨慎使用缓存 TTL 设太短基本没效果设太长结果又可能过期。5.2 通过火焰图定位性能卡点Flink 作业跑得不快先别急着骂集群。先看 Web UI 上每个算子的吞吐和延迟找到瓶颈算子再深入到 JVM 层面分析。火焰图是分析 CPU 性能最直观的手段。它能告诉你 CPU 时间到底花在哪个函数上是序列化太慢还是 GC 太频繁或者是某个外部调用阻塞了。Linux 下最常用的工具是 async-profiler下载后执行./profiler.sh -d 60 -e cpu -f /tmp/flink_flame.svg PID其中-d是采样时长-e cpu是采样 CPU 事件输出的/tmp/flink_flame.svg可以用浏览器打开。火焰图的横轴是采样占比纵轴是调用栈越宽的色块说明越热应该优先分析。拿到火焰图后重点看几类现象。如果看到org.apache.flink.runtime.io.network.buffer相关的调用栈很大往往是算子之间的数据传输有问题说明网络缓冲或序列化开销高。如果看到 GC 相关的栈很宽说明状态或对象分配太频繁可以调整堆内存、改用 RocksDB 状态后端或者优化代码减少对象创建。如果看到org.apache.kafka.clients.producer占用高说明 Kafka 写入或元数据拉取有问题。还有一种快速定位办法是用jstack连续抓几次线程栈for i in $(seq 1 30); do jstack PID jstack.log sleep 1 done然后用关键字统计每个线程栈出现的次数比如grep -A 20 flink jstack.log | sort | uniq -c | sort -nr。出现次数最多的栈对应的方法基本就是主要热点。这个方法虽然粗糙但在只能远程终端的情况下非常实用。给新手一个建议遇到性能问题先看指标再上工具不要凭感觉优化。Flink Web UI 的“Job Metrics”和“Task Metrics”面板里已经有很多关键指标先对比每个并行子任务的输入数据量、处理延迟、反压状态很多时候问题一眼就能看出来。6. 面试或复盘时躲不开的 Flink 问题6.1 高频问题Flink 和 Spark Streaming 有什么区别这个问题基本是流式计算岗位的必考题。可以从三个层面回答。从模型上Spark Streaming 传统上基于微批处理把数据切成小批量来做现在 Spark 3 虽然也有 Structured Streaming 的连续处理模式但架构核心仍然偏向批处理。Flink 则是真正的流处理引擎数据一到就处理延迟更低实时性更强。从状态容错上Flink 的分布式快照机制支持精确一次语义Spark Streaming 通过预写日志和状态更新也可以做到不丢不重但端到端的精确一次语义配合外部系统时Flink 的实现更成熟。从生态和编程模型上Spark 的优势在于完善的批处理和机器学习生态Flink 的优势在于流批一体和丰富的时间语义支持。如果业务核心是实时数据选 Flink 更合适如果是运维一个大的离线数仓偶尔跑实时任务Spark 也有它的位置。6.2 高频问题如何保证精确一次语义精确一次指的是每条数据对结果的影响只生效一次不重不丢。Flink 的机制包括两个关键部分。第一部分是分布式快照Checkpoint。Barrier 从 source 随数据流穿过所有算子算子收到 Barrier 后把状态快照保存到外部存储所有算子完成快照则代表一次检查点成功。作业故障后所有算子恢复到同一次检查点的状态source 也恢复到对应的位点这样已经产生的数据不会重复消费。第二部分是两阶段提交。Checkpoint 快照保证了状态一致但外部 Sink 也需要配合支持事务或幂等写入才能做到端到端不重不丢。Kafka Sink 支持事务提交配合 Flink 的 TwoPhaseCommitSinkFunction可以实现真正完整的精确一次语义。面试时如果能现场画出 Barrier 穿过算子的流程并且说清楚每个阶段要做什么面试官一般会比较满意。6.3 其他几个高频考点反压是另一个高频问题。上游算子产生数据的速度大于下游算子消费数据的速度压力会沿数据链路反向传播。排查方法是看 Web UI 上每个算子的“Backpressure”状态黄色或红色说明存在反压。解决思路一般是扩展并行度、优化瓶颈算子、调整缓冲区大小或者从根本上降低数据量。水位线和迟到数据也经常连着问。你可以这么说水位线用来推进事件时间迟到数据通过 allowedLateness 和侧输出流处理。allowedLateness 让窗口在触发后继续等待一段时间侧输出流把所有到得太晚的数据单独存放方便后续补算或排查。状态后端的问题也常见。HashMapStateBackend 适合状态小、要求低延迟的场景RocksDBStateBackend 适合状态大、超过堆内存的场景缺点是序列化和反序列化开销高一点。还有一个增量 Checkpoint 参数RocksDB 开启后性能会明显提升。6.4 学习路线建议如果你刚入门我的建议是先把官方文档“Concepts”部分过一遍然后动手部署本地集群跑通一个流式 WordCount再用 Flink SQL 做几张简单的实时报表。等基础打通后围绕状态、检查点、反压、背压这几个主题做深入实验最后再看源码。不要一开始就追新版本每个新版本都有新特性但对新手来说能把一个版本的生态吃透已经很难得。遇到报错先看日志和 Web UI再动手查资料。一步一步来不用着急。我个人的体会是Flink 入门最大的障碍不是语法而是思维切换。过去写批处理总想着“等数据齐了再算”写 Flink要习惯“数据来了就算状态自己记住”。等你想通了这一点再回头去看各种框架特性和面试题都会觉得顺理成章。这一章把概念底座打牢后面不管是做实时数仓、写连接器还是做性能调优都能少走很多弯路。
返回列表