ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Flink实时计算实战:构建新闻热搜分析系统与可视化大屏

2026/9/28 13:39:57 拓冰建站 浏览量
Flink实时计算实战:构建新闻热搜分析系统与可视化大屏 1. 项目定位与整体方案设计1.1 这项目到底在解决什么问题刚接触大数据实时计算的同学大概率会有这样一个困惑离线数仓那套Hive Spark SQL 定时调度已经能算不少东西了为什么还要搞一套 Flink 实时分析我自己的体会是当数据从晚上算昨天进化到秒级算现在的时候问题性质就完全变了。拿新闻热搜来说热搜榜的排序每分钟都在变一条突发新闻可能在几分钟内冲上榜首也可能半小时后就被淹没。如果还用 T1 的离线架构等你看清楚趋势的时候热点早就凉透了。这个实战项目本质上是构建一套采集 - 清洗 - 实时统计 - 结果存储 - 可视化的完整链路。数据源是新闻平台的热搜榜单或者新闻流通过 Flink 做实时流处理计算关键词热度、新闻上榜时长、涨幅趋势这些指标最终落到 MySQL 或者 Hive再配合 Flask ECharts 做一个可视化的热搜大屏。做完这套东西你不仅能熟悉 Flink 的核心 API还会把状态管理、窗口计算、检查点、连接器这些实战中绕不开的组件全部过一遍。适合谁来抄作业我觉得有三类人最合适。第一类是大数据专业的毕业生正愁毕业设计选题这个项目麻雀虽小五脏俱全从数据接入到可视化都有完整闭环第二类是准备大数据岗位面试的同学因为 Flink 的面试题翻来覆去就是窗口、状态、背压、容错这些点你亲手做过一遍和只看八股文的记忆深度完全不同第三类是工作中需要做实时大屏或者实时数仓的工程师可以参考里面的架构思路和数据流设计。我自己当时做这个项目前后花了大约三周用的是单机 Flink 集群 Kafka MySQL 的组合跑通之后又加了一些进阶功能下面把完整的思路和实践过程都写出来。1.2 为什么选 Flink 而不是 Spark Streaming选 Flink 做实时分析在当下的技术环境里基本不需要太多犹豫但为了照顾刚入门的朋友我还是把对比逻辑讲清楚。实时计算领域目前主要就是 Flink 和 Spark Streaming 两套体系Spark Streaming 的核心思路是微批次把流数据切成一小段一小段的小批量来处理延迟通常在秒级Flink 则是真正的事件驱动每条数据过来都能立刻触发计算延迟可以做到毫秒级。你可能会说热搜分析秒级够了呀为什么要毫秒级关键是 Flink 的窗口计算和状态管理比 Spark Streaming 成熟得多尤其是会话窗口、自定义窗口触发器这些场景写起来非常顺手。再有一个非常重要的原因是 Flink 的容错机制。实时任务最怕的就是程序挂了数据丢了Flink 的 Checkpoint 机制配合 Kafka 的偏移量管理能做到精确一次Exactly-Once的语义。这意味着哪怕是运行到一半集群重启任务也能从最近一次检查点恢复消息一条不丢、一条不重。做新闻热搜这种对数据准确性要求比较高的场景这个能力是刚需。我实际跑项目的时候故意 kill 掉过 TaskManager 来测试恢复效果重启之后数据能接上断点继续算这个体验在 Spark Streaming 里要费不少劲才能达到。架构上我的整体设计是一个六层结构数据采集层用爬虫或者第三方 API 拉取热搜数据投递到 Kafka 消息队列数据接入层用 Flink 自带的 Kafka Source 消费数据数据计算层做过滤、清洗、分词、窗口聚合结果存储层用 JDBC 连接器写入 MySQL同时落一份到 Hive 做历史归档最后是应用层用 Flask 提供接口ECharts 渲染大屏。这套架构最大的好处是每一层都可以独立替换比如你今天不想用 Kafka直接把 Source 换成自定义的 Data Source 也能跑明天想把存储从 MySQL 换成 ClickHouse只需要改 Sink计算逻辑完全不用动。2. 数据接入与预处理环节的实现细节2.1 自定义 Data Source 和 Kafka 接入怎么选项目里我见过不少同学一上来就写自定义 Data Source其实这是一个选择题。如果你是从零开始做 Demo想快速验证计算逻辑那么用 Flink 的addSource方法自定义一个 SourceFunction 是最简单的方案。大致逻辑是用 HttpClient 定时请求热搜接口解析返回的 JSON把数据封装成NewsHotEvent对象再用collect()发射到下游算子。这种方式的好处是依赖少、逻辑直观适合学习和调试。// 自定义 Data Source 的核心骨架 public class HotSearchSource extends RichSourceFunctionNewsHotEvent { private volatile boolean running true; Override public void run(SourceContextNewsHotEvent ctx) throws Exception { while (running) { // 1. 请求热搜接口 String response HttpUtil.get(https://api.example.com/hotsearch); // 2. 解析 JSON 为事件对象 ListNewsHotEvent events JSON.parseArray(response, NewsHotEvent.class); // 3. 发射数据 for (NewsHotEvent event : events) { ctx.collect(event); } // 4. 控制采集频率避免对接口造成压力 Thread.sleep(5000); } } Override public void cancel() { running false; } }但如果你想做一个更接近生产环境的项目Kafka 是绕不开的。热搜数据的特点是产生速度快、峰值波动大如果采集程序直接对接计算程序一旦 Flink 任务重启或者做检查点数据源就会积压或者丢失。引入 Kafka 相当于加了一个巨大的缓冲池采集端只管往 Kafka 里写Flink 只管从 Kafka 里读两边完全解耦。而且 Flink 的 Kafka Source 天然支持偏移量提交配合检查点能保证数据不丢不重这是自定义 Source 很难做到的。我个人的建议是学习阶段两种都做一遍先自定义 Source 把主流程跑通再切换到 Kafka 体验生产级的数据接入。2.2 数据清洗脏数据比你想的多得多很多人以为热搜接口返回的数据都是规范的真做了才知道里面的坑有多深。拿我做过的某个新闻平台热搜榜为例返回的字段大概有排名、标题、热度值、分类、发布时间等但实际拿到的数据经常出现这些问题标题里混入营销推广内容、热度值为负数、排名空缺、发布时间格式不统一、偶尔还有整条 JSON 解析失败。如果不加清洗直接进窗口计算结果基本不能看。我在项目里专门写了一个清洗函数挂在 Source 下游的第一个算子核心做了四件事。第一是格式校验过滤掉关键字段缺失或者类型不对的数据第二是去重用标题的哈希值做去重键配合 Flink 的 KeyedState 记录最近一段时间出现过的标题重复数据直接丢弃第三是热度值矫正对负数或者超出合理范围的值做截断处理第四是时间标准化把各种格式的时间统一转换成时间戳因为后面窗口计算完全依赖事件时间。这些逻辑看起来琐碎但在实时链路里非常重要因为离线任务错了还能回刷实时任务一旦脏数据流过去了结果错了很难追溯。这里要提示一个性能相关的细节清洗算子不要做太重的计算比如不要在这里调用外部接口做内容审核否则会成为整个链路的瓶颈。清洗逻辑应该保持无状态优先能不用 KeyedState 就不用因为状态越大检查点备份和恢复的成本越高。我当时为了去重引入了 RocksDB 状态后端后来发现单机跑的时候内存状态后端就够用了换回之后吞吐反而更高。2.3 分词与热词提取的实现思路新闻热搜分析除了统计榜单排名还有一个很常见的需求是提取热词。比如一整条新闻标题某地突发强降雨 多部门联动救援我们希望提取出强降雨救援这些关键词然后统计关键词的出现频次。这一步看起来简单实际操作起来有个容易踩坑的地方直接用现成的分词工具比如 HanLP、Jieba按默认模式分词往往会把突发强降雨切成突发强降雨统计出来的热词零散且没有意义。我的做法是两步走。第一步构建一个领域词典把新闻类的高频词汇灾害、救援、政策、科技、体育等以自定义词典的形式加载进分词器让分词结果偏向新闻语义第二步是做词性过滤只保留名词、动词、形容词去掉助词、介词、标点等无意义词元。处理完之后再交给 Flink 做窗口聚合统计每个时间窗口内关键词的排名。你可能会觉得这不就是离线分词吗但放在 Flink 里的区别是分词算子是有状态的词典可以从外部配置中心动态加载而不需要重启任务窗口统计也是实时的结果每五分钟刷新一次大屏上能看到热词排名的实时变化。代码上实现的核心是用flatMap做分词再用keyBy(word)window(TumblingProcessingTimeWindows.of(Time.minutes(5)))做滚动窗口计数。如果希望热词排名更平滑可以改成滑动窗口比如每 5 分钟输出一次最近 30 分钟的统计结果。滑动窗口能有效避免整点切分的抖动大屏展示出来的曲线会更自然。3. 核心计算逻辑与状态管理实战3.1 怎么设计热度值和涨幅指标实时热搜分析的指标设计决定了整个项目的价值。如果只是简单地把接口返回的热度值原样展示那不叫实时分析叫数据搬运。项目中我做了三个核心指标当前热度值、热度涨幅、上榜时长。当前热度值就是接口返回的原始值但要做归一化处理因为不同平台的量纲差别很大有的平台是几十万有的平台只有几千不归一化没法横向对比。我的归一化公式很简单score (value - min) / (max - min)min 和 max 取自当前窗口内的历史极值这个极值通过 Flink 的累加器状态来维护。热度涨幅是更有意思的一个指标。它衡量的是某条新闻在相邻两个窗口之间的热度变化率公式是(currentScore - lastScore) / lastScore。涨幅超过一定阈值比如 50%的新闻说明正在快速发酵需要在可视化大屏上标红置顶。这个计算需要跨窗口访问上一步的结果实现方式是在内存中维护一个 KeyedStatekey 是新闻标题value 是上一次窗口的热度值。这里要注意状态的无界增长问题因为新闻数量是无限的如果不加清理机制状态会越来越大最终拖垮任务。我的方案是给状态设置 TTL比如 24 小时没有更新的新闻自动过期清除。最后一个指标是上榜时长统计一条新闻从首次出现到当前时间持续了多久。这个用 Flink 的ValueState记录首次时间戳每条数据达到时计算差值即可。两个状态首次时间、上次热度值都挂在新闻标题这个 key 下状态后端选择 RocksDB一方面是为了支持大状态另一方面是后续如果要扩展到百万级新闻标题也不至于内存爆炸。3.2 窗口类型选型滚动、滑动还是会话窗口是 Flink 实时计算的核心概念很多新手学的时候觉得简单真到业务场景就不知道怎么选。拿新闻热搜来说三种窗口都有各自的适用场景。滚动窗口Tumbling Window把数据流切成固定大小的段比如每 5 分钟一个窗口窗口之间没有重叠适合做周期性的快照统计比如展示当前五分钟热搜榜滑动窗口Sliding Window则允许窗口之间有重叠比如每 1 分钟滑动一次、窗口长度 10 分钟这样每隔一分钟就能产出一次覆盖过去十分钟的统计结果适合做趋势平滑展示会话窗口Session Window是根据事件之间的间隔动态划分窗口超过指定间隔没有新数据就关闭当前窗口适合分析用户行为序列或者突发事件的完整生命周期。我在这项目里混合使用了两种窗口。实时榜单用的是 10 分钟滑动、5 分钟滑动的设计这样热搜排名的变化会非常平滑不会有每五分钟跳变一次的生硬感热词趋势分析用的是滚动窗口因为热词统计的结果要写入数据库做后续的离线分析滚动窗口的语义更干净容易对齐时间分区。如果你的业务是监控突发新闻从发酵到消退的完整过程那会话窗口会更合适比如设置 30 分钟无更新就认为这条新闻已经不再热门。有一个细节值得单独提一下处理时间和事件时间的区别。如果你用ProcessingTime窗口的划分取决于 Flink 任务所在机器的系统时间优点是延迟低、不需要等待乱序数据缺点是结果不准确一旦上游数据延迟到达会被划分到错误的窗口。我在项目里最终选择了EventTime Watermark因为热搜数据在 Kafka 里传输可能会有延迟只有按事件发生的时间来算才能保证窗口结果的正确性。Watermark 的设置我当时取了 5 秒的延迟允许最多 5 秒内的乱序数据超过这个范围的就丢弃。3.3 自定义 SinkJDBC 写入 MySQL 和 Hive计算完成之后的结果需要写到外部存储供可视化层读取。我在项目里写了两类 Sink一个写入 MySQL 供实时大屏查询另一个写入 Hive 做历史归档。先讲 MySQL 写入这里强烈建议你用JdbcSink而不是自己在RichSinkFunction里拼 JDBC因为前者内置了批量提交和重试机制。批量提交的时机很关键我设置的是攒够 1000 条或者间隔 10 秒就批量 flush 一次这样能大幅减少数据库连接的开销。不过这里有一个坑我必须提醒你Flink 官方提供的JdbcSink在某些版本下对字段映射的处理方式有变化尤其是 SQL 语句中的占位符顺序和 POJO 字段顺序不一致的时候容易报字段类型不匹配的错误。我踩过一次比较恶心的坑是MySQL 表里的时间字段是 datetime但 Flink 这边的时间戳是 Long 类型直接写入会抛异常需要在写入前把 Long 转换为 Timestamp或者让 POJO 里的时间字段直接定义成java.sql.Timestamp。// 使用 JdbcSink 批量写入 MySQL DataStreamHotSearchResult resultStream ...; resultStream.addSink(JdbcSink.sink( INSERT INTO hot_search_result(title, score, increase_rate, window_start, window_end) VALUES (?,?,?,?,?) ON DUPLICATE KEY UPDATE scoreVALUES(score), (ps, event) - { ps.setString(1, event.getTitle()); ps.setDouble(2, event.getScore()); ps.setDouble(3, event.getIncreaseRate()); ps.setTimestamp(4, event.getWindowStart()); ps.setTimestamp(5, event.getWindowEnd()); }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchInterval(Time.seconds(10)) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(jdbc:mysql://localhost:3306/hotsearch) .withDriverName(com.mysql.cj.jdbc.Driver) .withUsername(root) .withPassword(password) .build() ));写入 Hive 则是另一套思路。Hive 本身不擅长高频写入更常见的方式是用 Flink 的 StreamingFileSink 先写成 Parquet/ORC 文件再通过 Hive 的分区表来读取。我在项目里是按照dt字段做分区每个小时滚动一个分区这样离线分析的时候直接按分区扫描效率很高。这里踩过的坑是Flink 写 Hive 需要引入额外的 Hive 依赖而且 Flink 版本和 Hive 版本必须匹配我用 Flink 1.14 配 Hive 3.1.2 的时候光是调 HiveCatalog 的兼容性就花了一个下午。如果你只是做项目演示我建议先把结果写到 MySQLHive 作为可选项后期再接入避免一开始就被环境问题劝退。3.4 状态后端选择内存还是 RocksDB状态后端的选型直接决定了任务能撑多大的数据量。Flink 默认的内存状态后端HashMapStateBackend把状态放在 JVM 堆内存里读写速度极快适合状态量小的场景但受限于堆内存大小一旦状态超过几 GB 就会频繁 GC 甚至 OOM。RocksDB 则把状态存储在本地磁盘堆外内存开销小能支撑 T B 级别的状态但代价是读写性能比纯内存模式低一个数量级。我的经验准则是如果状态量在 1GB 以内无脑用内存状态后端简单高效如果状态量可能无界增长或者你预期未来要横向扩展那就直接用 RocksDB。前面提到的新闻去重状态和涨幅计算状态在单机演示模式下很小内存完全够用但我仍然换了 RocksDB原因有两个第一是 RocksDB 支持增量检查点做 Checkpoint 的时间更短第二是想提前体验生产环境的配置方式毕竟面试官大概率会问状态后端怎么选这个问题。还有一个和状态强相关的配置是 Checkpoint。检查点间隔我设置为 60 秒模式设置为 EXACTLY_ONCE超时时间 30 秒最大并发检查点数为 1。如果你做的场景可以接受少量重复数据也可以把模式改为 AT_LEAST_ONCE检查点间隔可以缩短到 30 秒恢复速度会更快。这里要说一句检查点间隔不是越短越好间隔太短会导致频繁触发全量快照反而影响正常处理性能。4. 部署调优与问题排查实录4.1 集群部署策略单机、Standalone 还是 Flink on YARN项目跑通之后你肯定会考虑部署到集群上。Flink 的部署方式主要有三种单机本地模式、Standalone 集群、Flink on YARN。单机模式适合开发和调试你不用管资源调度一条命令就能起一个 Flink 集群我在项目初期基本都用这个模式配合 IDEA 的本地调试非常方便。Standalone 集群是独立的 Flink 集群需要自己管理 TaskManager 的资源分配适合小规模场景但生产环境其实很少这么用了因为资源和任务之间没有自动隔离。如果你是在 Hadoop 生态圈里做项目那最佳实践是 Flink on YARN。把 Flink 任务提交到 YARN 上YARN 负责分配和管理容器资源任务挂掉之后会自动拉起资源利用率也高得多。部署策略上有一个很容易被忽视的问题TaskManager 的并行度和内存配置。很多人图省事直接用默认配置结果任务跑起来之后发现 CPU 利用率很低或者频繁 GC。我的建议是先把并行度设为集群可用核心数的一半观察处理延迟和吞吐量再逐步调整。本地调试和集群部署还有一个区别本地模式下 Kafka 的地址可以直接用 localhost:9092但部署到集群之后必须改为集群内可达的主机名或者负载均衡地址。我遇到过的真实案例是任务部署到 YARN 上之后一直报 Kafka 连接超时排查了半天才发现是/etc/hosts里没有配置 Kafka 的机器名映射添加配置后立即恢复。这些小问题在实战中非常消耗时间提前检查环境配置能省掉很多麻烦。4.2 Flink CDC Pipeline 与数据血缘进阶玩法项目做完基础版本之后如果你还有精力我强烈建议研究一下 Flink CDC Pipeline。这个功能可以从数据库的 binlog 中实时捕获变更数据然后通过 Pipeline 直接同步到目标端。对应到这个项目里一个很自然的应用场景是MySQL 中存储的热搜结果表有更新比如热度值变化通过 CDC 实时同步到 Elasticsearch 或者 Hive供下游的搜索和分析使用。相比传统的批量同步CDC 的延迟能降低到秒级而且对源库的影响非常小。这里要坦白说一句Flink CDC 的安装部署比普通 Flink 任务要麻烦一些主要是版本兼容性问题。我在项目中尝试过把 MySQL 的变更数据通过 Flink CDC Pipeline 同步到另一个 MySQL 实例过程中遇到了不少问题包括需要提前下载对应的连接器依赖、确认 MySQL 的 binlog 格式为 ROW、设置合理的并行度避免 Source 端压力过大。这些问题排查下来我对 Flink 的运行时机制理解又深了一层。另外还有一个值得关注的功能是数据血缘。用 OpenMetadata 这样的元数据管理工具可以自动采集 Flink 任务的血缘关系也就是搞清楚一张结果表的数据是从哪个 Source 表、经过哪些算子计算出来的。做大数据治理的时候血缘关系是回答这个数据怎么来的这个问题的关键。我一开始觉得血缘关系只是锦上添花直到有一次需要排查一个异常指标回溯数据链路的时候发现中间有个算子没有按预期过滤血缘图帮我快速定位了问题。如果你的项目最终要交付给团队使用加上元数据采集是很有价值的。4.3 性能调优从火焰图出发当你的项目要处理的数据量变大之后性能问题就会暴露出来。一个非常实用的调优工具是 JFRJava Flight Recorder配合火焰图分析。火焰图能直观展示 CPU 时间都花在了哪些方法上是定位 CPU 密集型问题和锁竞争问题的利器。我跑项目的时候发现开启火焰图之后可以清楚看到某些序列化和反序列化的操作占据了大量 CPU这提示我可以调整数据类型来减少序列化开销。Flink 里面有一个常见的序列化陷阱如果用自定义的 POJO 类型尽量使用 Flink 自带的 TypeInformation 而不是 Java 原生序列化否则性能会差好几个数量级。我观察到一个非常典型的实例同样的逻辑用 POJO 类型跑吞吐量是使用GenericType的 3 倍左右。所以定义数据模型的时候最好直接用简单的case classScala或者带无参构造方法的普通 Java 类并且字段类型定义明确不要用Object。内存调优方面有个参数特别值得注意taskmanager.memory.process.size。Flink 1.14 之后内存模型变得复杂分为堆内存、托管内存用于 RocksDB 等、直接内存等几个部分。如果内存配置不合理任务会频繁出现ContainerMemoryExceeded异常我在部署到 YARN 上时也遇到过这个问题。最简单的处理方法是预留总内存的 15%-20% 给系统开销托管内存根据是否使用 RocksDB 来调整如果用了 RocksDB托管内存建议设置为总内存的 30%-50%。4.4 常见问题速查表做完整项目之后我把遇到的高频问题整理成了一个查表每次遇到类似情况能快速定位。下面几张表是我实际踩坑的总结希望能帮你省一些没有必要的排查时间。问题现象可能原因解决方案Kafka Source 一直报连接超时集群环境未配置主机名映射检查/etc/hosts确认 Kafka 地址可达JDBC Sink 写入 MySQL 报字段类型错误Long 类型与 datetime 字段不匹配写入前把 Long 转为 Timestamp写入 Hive 后查询不出数据分区字段类型不一致或未执行msck repair统一分区字段类型刷新 Hive 分区元数据窗口结果不更新使用 ProcessingTime 但数据乱序严重切换为 EventTime WatermarkCheckpoint 一直失败状态后端存储空间不足或并行度过高调整状态后端配置降低并行度任务恢复后数据重复Checkpoint 模式为 AT_LEAST_ONCE改为 EXACTLY_ONCE注意下游幂等性Flink 面试题里最常问的几个问题在这个项目里都能找到具体的答案Exactly-Once 是怎么实现的两阶段提交 检查点背压是怎么产生的下游处理速度跟不上上游数据量以及怎么缓解调整并行度、优化算子逻辑Watermark 和乱序数据的处理策略是啥延迟触发窗口、侧输出迟到数据。纸上谈兵很容易但只有当你真正在日志里看到背压报警、看到数据倾斜导致某个子任务积压几百万条数据的时候才会对这些机制有刻骨铭心的理解。另一个容易被忽略的点是 Flink 任务的全局配置比如ExecutionConfig里开不开对象重用enableObjectReuse。开启之后可以减少对象创建和 GC 压力但代价是同一个对象会在多个算子间复用如果你在业务代码里修改了对象内容可能会产生副作用。我的建议是如果你不熟悉内部机制保持默认关闭先求稳再求快。4.5 环境搭建与踩坑记录最后再单独说一句环境搭建。网上关于 Flink 的安装配置教程特别多头歌平台上也有一系列Flink 安装配置到部署的关卡但我发现很多人卡在同一个地方Flink 和 Java、Scala 的版本兼容关系。Flink 1.14 需要 Java 8 或 Java 11如果你用的是 Java 17 或者更高版本大概率会碰到UnsupportedClassVersionError。还有个常见问题是本地能跑通、提交到集群就报错这通常是因为本地依赖和集群依赖冲突比如把flink-shaded-hadoop打进了 fat jar。我的经验是提交到 YARN 的任务依赖尽可能用 provided 模式让集群提供公共的 Flink 和 Hadoop 组件只把项目自身的业务代码打进去。启动参数上开发环境和生产环境也需要区分。本地调试时我习惯直接指定-m localhost:8081让任务跑在本地方便看 Web UI集群部署时则用-t yarn-session或者-t yarn-per-job的方式提交。这里多说一句如果每条新闻的首次出现时间状态数据丢失会导致上榜时长从 0 重新计算大屏上会出现时长回退的诡异现象。这个现象背后的原因是状态没有做持久化或者恢复失败。我在排查中发现如果检查点路径配置错误任务重启之后状态是干净的之前的累计数据全部丢失。所以检查点路径最好配置在分布式文件系统上比如 HDFS而不是本地磁盘否则换个节点运行就找不到之前的检查点。检查点对实时任务来说就是命根子再怎么强调都不为过。我这套项目做完之后复盘时最大的体会是实时计算任务的难点不在于写几个算子而在于把整个链路的稳定性和可维护性撑起来。数据怎么不丢、状态怎么恢复、结果怎么对账、性能怎么调优这些才是区分有没有生产经验的分水岭。建议你拿到项目之后先照着完整跑一遍然后故意去破坏一些环节杀掉 TaskManager、停掉 Kafka、写入脏数据观察 Flink 的反应这个过程比看十篇教程都有价值。等你把这套流程走熟了再去看 Flink 源码或者官方文档很多以前看不懂的设计理念会突然就通了。