跳转至

Flink

Flink是一个分布式的计算框架,主要用于对有界和无界数据流进行有状态计算。 有界数据流就是指离线数据,有明确的开始和结束时间,无界数据流就是指实时数据,源源不断没有界限。 有状态计算指的是在进行当前数据计算的时候,我们可以使用之前数据计算的结果。 Flink还有一个优点 就是集成了很多高级的API,比如DataSet API,DataStream API,Table API 和FlinkSQL。

  • Framework Heap(框架内存)
    • Flink 框架本身占用的内存
    • 堆内:默认128M
    • 堆外:默认128M,以直接内存分配
  • Task Heap(Task内存)
    • 用于存放 Flink 应用的算子及用户代码
    • 堆内:无默认值,会用Flink总内存减去其他部分内存得出
    • 堆外:默认0M,以直接内存分配,一般不用
  • Managed Memory(托管内存)
    • 用于存放中间结果缓存、排序、哈希表等,以及 RocksDB 状态后端。
    • 纯堆外内存,由 MemoryManager 管理
    • 通过参数可以调节 托管内存占 Flink 总内存的比例,默认0.4
  • Network Memory(网络缓存)
    • 在 Task 与 Task 之间进行数据交换时(shuffle),需要将数据缓存下来
    • 默认最小64M,最大1G
    • 网络缓存占 Flink 总内存 的比例,默认值 0.1
  1. 第一,计算速度的不同。Flink是真正的实时计算框架,而Spark Streaming 是一个准实时微批次的计算框架。也就是说,Spark Streaming的实时性比起Flink,差了一大截。
  2. 第二,架构模型的不同。Spark Streaming 在运行时的主要角色包括:Driver、Executor,而 Flink 在运行时主要包含:Jobmanager、Taskmanager
  3. 第三,时间机制的不同。Spark Streaming 只支持 "处理时间" 语义,而Flink支持的时间语义包括 "处理时间"、"事件时间"、"注入时间" ,并且还提供了watermark(水位线) 机制来处理迟到数据。

Flink的shuffle

Flink 的 Shuffle 是指数据在任务间传输的核心机制,直接影响作业性能(吞吐、延迟)和资源利用率。设计Flink shuffle的核心目标是实现低延迟、高吞吐、高容错的高效传输

  • 核心组件
    • TaskManager(TM) 每个 TM 包含多个 Task Slot,负责执行具体任务(如 map()keyBy())。
    • 网络栈(Netty) 基于 Netty 4 实现异步非阻塞通信,支持零拷贝(Zero-Copy)技术减少内存拷贝。
    • 缓冲区(Buffers)
      • 内存缓冲区:优先使用堆外内存(Direct Memory),避免 JVM 堆压力。
      • 文件缓冲区:内存不足时溢出到磁盘,平衡吞吐与内存使用。
  • shuffle的传输方式
    • 点对点传输(Push-Based) 上游任务直接将数据发送给下游目标 Task,无需集中式 Shuffle Manager,避免单点瓶颈。
    • 动态分区(Dynamic Partitioning) 根据 Key 的分布动态调整分区策略,减少数据倾斜(如热点 Key)。
    • 合并请求(Request Batching) 下游 Task 合并多个数据请求,减少网络往返次数。
  • 流处理 vs. 批处理优化
    • 流处理
      • 低延迟:采用 无界流分区,实时传输数据。
      • 反压机制:通过缓冲区水位(High/Low Watermark)动态调整发送速率。
    • 批处理
      • 高吞吐:使用 有界流分区,批量发送数据以减少序列化开销。
      • 排序优化:对 Key 进行预排序,提升下游聚合效率。
  • 容错与状态恢复
    • Barrier 机制 在 Checkpoint 时插入 Barrier,确保上下游状态一致性。
    • Shuffle 数据重放 故障恢复时,从最新 Checkpoint 重新计算并传输丢失数据。
  • 性能优化策略
    • 序列化优化 使用 Kryo/Avro/Protobuf 等高效序列化框架,减少 CPU 开销。
    • 压缩传输 对高频小数据启用压缩(如 Snappy),降低网络负载。
    • 避免数据倾斜 通过 预聚合 或 二次 Key 设计 分散计算压力。

假设一个 keyBy() 操作后接 window()

  1. 上游 Task 根据 Key 将数据发送到对应下游 Task。
  2. 下游 Task 缓冲数据,触发计算后输出结果。
  3. Checkpoint 时,Barrier 确保上下游状态同步,故障时可从 Savepoint 恢复。

以单作业模式为例

  1. 作业提交。用户提交作业,指定YARN作为资源管理器。Flink 客户端将作业 JAR 包和配置信息发送给 YARN 的 ResourceManager(RM)
  2. 启动 ApplicationMasterYARN 调度器分配一个 Container 并选择合适的NodeManager(NM)启动 Flink 的 ApplicationMaster(AM)。AM 负责与 YARN 交互,申请资源并管理 Flink 集群生命周期。
  3. 初始化 JobManager AM 向 ResourceManager 申请资源,启动 JobManager 的 Container。JobManager 负责作业调度、任务划分和状态管理。
  4. 资源协商与 TaskManager 启动 JobManager 根据作业需求(如并行度)向 ResourceManager 申请更多 Container。AM 将这些 Container 分配给 TaskManager(TM),TM 启动后向 JobManager 注册。
  5. 任务执行 JobManager 将作业拆分为多个 Task,分配给空闲的 TM 执行。TM 通过 Slot 资源执行任务,并定期向 JobManager 汇报状态。
  6. 作业完成与资源释放 作业执行成功后,AM 向 ResourceManager 释放所有 Container,关闭集群(单作业模式)。若启用 Checkpoint,状态可持久化到外部存储。

以下是一个最简单的 Apache Flink 数据处理任务示例,使用 Java 实现,这个程序从一个整数集合读取数据,将每个数字乘以 2,然后输出结果

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class SimpleFlinkJob {
    public static void main(String[] args) throws Exception {
        // 1. 创建执行环境
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 2. 创建数据源(从集合中读取数据)
        DataStream<Integer> numbers = env.fromElements(1, 2, 3, 4, 5);

        // 3. 数据处理:将每个数字乘以2
        DataStream<Integer> result = numbers.map(number -> number * 2);

        // 4. 输出结果(打印到控制台)
        result.print();

        // 5. 执行任务
        env.execute("Simple Flink Job");
    }
}

Flink 通过 Checkpoint + 可靠 Source/Sink + 事务机制,在故障时从持久化状态恢复,确保无数据丢失(Exactly-Once)

  • Checkpoint 机制
    env.enableCheckpointing(60000); // 每分钟触发一次
    env.getCheckpointConfig().setCheckpointStorage("hdfs://path");
    
  • Source 可靠性
    • 偏移量管理:
    • Kafka Source:消费者仅在 Checkpoint 完成后提交偏移量(enable.auto.commit=false),确保故障时从已提交的偏移量重新消费。
    • 其他 Source:需支持重放机制(如文件 Source 记录读取位置)。
    • 状态后端
  • Sink 可靠性
    • 事务性 Sink:例如Kafka Producer,通过启用幂等性(enable.idempotence=true)和事务(transactional.id)**,确保数据不重复且可靠提交。
    • 两阶段提交(2PC):Sink 在 Checkpoint 完成前预提交数据,确认后正式提交,避免数据丢失。
    • 端到端 Exactly-Once:Source、Flink 处理逻辑、Sink 均支持事务,保证全局数据一致性
  • 背压机制
    • 当下游处理能力不足时,通过背压(Backpressure)通知上游降低发送速率,避免内存溢出导致数据丢失。

首先了解一下什么是一致性级别: 1. at-most-once:什么都不干,既不恢复丢失的状态 ,也不重播丢失的数据 2. at-least-once:一些事件可能被处理多次 3. exactly-once:没有事件丢失,并且对于每个事件,有且仅有处理一次

端到端的一致性实现保证包括: - 内部依赖checkpoint来保证 - source端需要source支持可重置偏移量,例如主要使用的kafka,如果要实现精准一次,作为source可以重发,然后由Flink维护偏移量,作为状态存储,指定隔离级别为 READ_COMMITTED 才不会消费到未正式提交的数据,避免重复消费;kafka作为sink端的时候,官方的推荐实现是基于2pc,需要指定sink的语义为Exactly-once,能保证写入的精确一次。 - sink端需要保证从故障恢复时,数据不会重复写入外部系统,主要的方法有两种:幂等写入和事务性写入 - 幂等写入就是同一份数据无论写入多少次都只保留一份结果 - 事务性写入有两种实现方式:WAL和2PC - WAL也叫预写日志,是先把结果数据写入log文件中,然后在收到checkpoint完成的通知的时候一次性写入外部系统 - 2PC也叫两阶段提交,预提交阶段,sink任务会启动一个事务,对于每个checkpoint,会将接下来所有接收的数据添加到事务里,然后 将这些数据写入外部sink系统,但不会提交它们,只有当收到checkpoint 完成的通知时,它才会正式提交事务,实现结果的真正写入。

整个系统的端到端的一致性级别取决于所有组件中一致性最弱的组件

  • 如果下级存储不支持事务,例如具体实现是幂等写入,那就需要下级存储也具有幂等性写入的特性。比如结合 HBase 的 rowkey 的唯一性、数据的多版本,实现幂等,或者结合 Clickhouse 的 ReplicingMergeTree 实现去重,查询时加上 final 保证查询一致性
  • 实际项目中,完全的精准一次为影响数据可见性,要等到第二次提交下游才能消费到。在数仓项目中,我们的使用时只考虑了至少一次,Kafka Source 设置的隔离级别是 READ_UNCOMMITTED ,Kafka Sink 也没有使用 exactly-once 来保证时效性。

简单介绍一下Checkpoint机制

Checkpoint检查点其实就是所有任务的状态在某个时间点的一份快照,且这个时间点,必须是所有任务都恰好处理完一个相同的输入数据的时候。

  • Checkpoint 的特性
  • 定期触发:Flink 定时(如每分钟)或定量(如每 1000 条数据)触发 Checkpoint,将状态数据(如算子状态、窗口计数)持久化到外部存储(如 HDFS、S3)。
  • Barrier 对齐:Checkpoint 时插入 Barrier,确保所有并行任务的状态同步,避免数据不一致。
  • 故障恢复:作业失败时,从最近完成的 Checkpoint 恢复状态,重新消费未确认的数据。

  • Checkpoint 的步骤

  • Flink应用在启动的时候,Flink的JobManager创建了CheckpointCoordinator
  • CheckpointCoordinator(检查点协调器) 周期性地向该流应用的所有 source 算子发送 barrier
  • 当某个source算子收到一个barrier时,便暂停数据处理的过程,然后将自己的当前状态制作成快照,并保存到指定的持久化存储(hdfs) 中,最后向CheckpointCoordinator 报告自己快照的制作情况,同时向自身所有下游算子广播该barrier,然后恢复数据处理
  • 下游算子收到barrier之后,会暂停自己的数据处理过程,然后将自身的相关状态制作成快照,并保存到指定的持久化存储中,最后向 CheckpointCoordinator 报告自身快照情况,同时向自身所有下游算子广播该barrier,恢复数据处理。
  • 每个算子按照上面这个操作不断制作快照并向下游广播,直到最后barrier传递到sink 算子,快照制作完成
  • 当CheckpointCoordinator 收到所有算子的报告之后,认为该周期的快照制作成功; 反之,如果在规定的时间内没有收到所有算子的报告, 则认为本周期快照制作失败。
  • Checkpoint机制是Flink保证数据容错性的核心机制,它可以将正在处理的状态保存在持久化存储中,以便在出现故障时能够恢复数据并继续进行处理。
  • 保存点(savePoint)是手动触发的全局一致性状态快照,主要用于版本升级作业重新部署时的恢复。
特性 Checkpoint(检查点) Savepoint(保存点)
触发方式 自动(定时/定量) 手动(用户主动触发)
确保 Exactly-Once 语义 全局一致性状态快照,主要用于版本升级或作业重新部署时的恢复
生命周期 与作业绑定,可能自动清理 独立存在,长期保存
并行度调整 通常不支持(需配置) 支持
典型场景 故障恢复、日常容错 作业更新、迁移、调整拓扑
跨作业恢复 仅恢复当前作业,重点在于自动容错 可用于不同作业(需兼容),多用于故障恢复

checkpoint 和 barrier 是同时进行的吗?(checkpoint 中barrier的两种对齐方式)

对齐主要针对多流算子的处理逻辑 - 如果是对齐检查点,则不会同时进行,barrier 对齐确保了所有输入流中的 barrier 都到达该算子后,才能进行 checkpoint。在数据强一致性且流量平稳的场景则使用默认的对齐。 - 如果是非对齐检查点,则会同时进行,因为每个算子的 barrier 都是独立触发的,不会等待其他算子的 barrier 到达,所有未处理的数据都被打包进了状态里,恢复时可以直接从 Barrier 位置继续消费。在低延迟、反压的场景首选。

非Barrier对齐可以保证精准一致性吗

  • 即使跳过Barrier对齐,Exactly-Once语义仍可保证,原因如下:
  • 状态恢复的完整性 非对齐检查点会将Barrier之后未处理的数据一并保存到快照中。恢复时,系统会从快照中重新处理这些数据,确保没有数据丢失或重复。
  • 端到端一致性支持 若Sink端支持两阶段提交(如Kafka事务写入),即使跳过对齐,整个处理链路仍可保证端到端的Exactly-Once。

Flink其中一个task的checkpoint失败后,算作整个checkpoint失败吗

不会,Flink的checkpoint机制是为每个任务独立进行的

  1. processing time(处理时间):表示执行算子的本地系统时间
  2. event time(事件时间):表示事件创建的时间,通常由事件中的时间戳描述
  3. ingestion time(注入时间):表示数据进入 Flink 的时间

  4. 在Flink的流式处理中,绝大部分的业务都会使用event Time

  1. 使用 watermark 设置延迟时间
  2. windowallowedLateness 方法,可以设置窗口允许处理迟到数据的时间
  3. 在Watermark触发窗口计算后,窗口不会立即销毁
  4. 额外保留一段时间继续接收迟到数据并重新触发计算
  5. windowsideOutputLateData 方法,可以将迟到的数据写入侧输出流
[数据流]

├─ Watermark(延迟2分钟)  触发窗口计算
   
   ├─ 窗口内数据  正常计算
   
   └─ 迟到2分钟内数据  纳入窗口

├─ allowedLateness(3分钟)  窗口保留期
   
   ├─ 迟到3分钟内数据  重新触发计算
   
   └─ 超过3分钟  侧输出流

└─ sideOutputLateData  窗口销毁,进入侧输出流,人工处理/特殊统计

数据乱序:只有在考虑事件时间语义的情况下,才会发生乱序(到达窗口的事件先后顺序和事件时间先后顺序 不一致)的情况。

watermark本质是一种特殊的时间戳,即watermark = 当前事件发生的最大事件时间- 设定的延迟时间,作用就是为了让事件时间慢一点,等迟到的数据都到了,才触发窗口计算,目的是为了让迟到的 窗口条件内 的数据 可以被纳入 本窗口的计算,从而降低错误率。 例如项目中开了一个5秒的窗口,但是2秒的数据在5秒数据之后到来, 那么5秒的数据来了,是否要关闭窗口呢?可想而知,关了的话,2秒的数据就丢失了,如果不关的话,我们应该等多久呢?所以需要有一个机制来保证可以在一个特定的时间后,关闭窗口,这个机制就是watermark。 所以,当watermark等于窗口时间的时候,就会触发窗口关闭并计算

状态可以理解为一个本地变量

Flink中主要有两种类型的状态,包括

  • operator state(算子状态),和 key 无关。本类型的状态,对于同一任务而言,是共享的;
  • keyed state(键控状态), 和 key 相关。本类型的数据,每一个key都会保存一个状态

状态后端就是用来进行状态的 存储、访问和维护 的东西,分为三种:(持久化存储的三种方式)

  • EmbeddedRocksDBStateBackend,存在本地磁盘,靠RocksDB管理,读写有序列化开销,适合大状态生产
  • HashMapStateBackend,将状态存在TaskManager的JVM堆内存里面,读写快但受到堆的大小限制,不支持增量检查点。是项目所使用的状态后端

状态编程去重时是如何选择状态后端

核心看数据量和去重窗口大小 1. HashMapStateBackend:适合中小规模数据,内存读写快,延迟低,搭配FSCheckpointStorage存检查点也能保证容错 2. EmbeddedRocksDBStateBackend:适合大规模数据或者长窗口去重,内存读写慢,延迟高,搭配因为磁盘存储可以抗住TB级状态,增量检查点还能减少存储开发,避免内存溢出。

--- 下面的状态后端已被废弃,不建议使用,仅保留历史记录 1. 内存状态后端(MemoryStateBackend) 1. 特点:将状态数据存储在Java堆内存中,读写速度快,适用于状态数据量小、对性能要求极高的场景。 2. 使用场景:在进行实时去重时,如果去重的数据集较小,例如在一个短时间窗口内对少量数据进行去重,使用内存状态后端可以获得很好的性能。 2. 堆外内存状态后端(OffHeapMemoryStateBackend) 1. 特点:把状态数据存储在堆外内存,能避免Java堆内存溢出问题,同时可以利用操作系统的内存管理机制,在一定程度上提高内存的使用效率。 2. 使用场景:当去重操作涉及的数据量较大,但又希望有较好的性能时,堆外内存状态后端是一个不错的选择。例如,对中等规模的实时数据流进行去重,且对内存管理有较高要求时可使用。 3. RocksDB状态后端(RocksDBStateBackend) 1. 特点:基于RocksDB存储状态数据,支持大规模状态存储,具有持久化和容错能力。即使在Flink作业失败或重启时,状态数据也能恢复。 2. 使用场景:对于大规模的实时数据去重,尤其是需要处理长时间窗口或大量历史数据的去重场景,RocksDB状态后端能更好地应对,它可以将大量状态数据持久化到磁盘,避免内存不足的问题。


Connect 算子和 Union 算子的区别

  1. 第一点,union算子的两个流类型必须是一样的,而connect算子的两个流类型可以不一样,因为connect之后,内部其实两个流还是独立,一般还需要使用comap来进行转换
  2. 第二点,union算子可以连接多个流,而connect算子只能连接两个流

interval join介绍

  • interval join只支持事件时间的场景
  • 只能支持两条流的关联
  • 在右流上划分一个范围区间,左流关联右流
DataStream API SQL API
编程方式 是一种基于 Java、Scala 等编程语言的面向对象的 API,开发者可以通过编写代码来实现各种复杂的业务逻辑,对数据的处理过程有更精细的控制 通过 SQL 语句来定义数据处理逻辑,更适合熟悉 SQL 的开发者,它提供了一种声明式的方式来处理数据,将数据处理的逻辑描述与具体的实现分离。
适用场景 适用于处理复杂的实时流计算场景,如需要对数据进行多步转换、聚合、窗口操作以及与外部系统进行深度集成等情况。例如,在实时监控系统中,需要对传感器数据进行实时分析,检测异常行为,可能涉及到复杂的规则匹配和状态管理,DataStream API 能更好地满足这些需求。 常用于对结构化数据进行简单的查询、过滤、聚合等操作,尤其适用于批处理和交互式查询场景。比如,在数据分析场景中,分析师需要快速对大量的历史数据进行统计分析,使用 SQL API 可以方便地编写查询语句来获取所需的结果
性能优化 开发者可以通过自定义函数、优化数据结构等方式进行深度的性能优化,以满足特定的性能要求。 Flink 会自动对 SQL 语句进行优化,如查询计划的生成、数据的分区裁剪等,但在某些复杂场景下,可能无法达到与 DataStream API 相同的性能优化程度。
代码可读性和维护性 代码相对复杂,尤其是在处理复杂业务逻辑时,可能会涉及到大量的代码编写和调试,但对于熟悉面向对象编程的开发者来说,代码的结构和逻辑比较清晰,易于维护。 SQL 语句具有较高的可读性,对于熟悉 SQL 的人来说,很容易理解查询的目的和逻辑。但如果业务逻辑复杂,可能需要编写复杂的 SQL 语句,此时维护性可能会受到一定影响。

这个问题我从两个方面说下自己的理解吧,时效性和准确性,然后分别聊一下 DataStream API 和 SQL API 是怎么做到的。

先说时效性。 Flink 的核心设计就是基于流处理的,所以不管是 DataStream 还是 SQL,底层都是同一个流执行引擎,天生就是来一条处理一条,不用等批攒够了再跑,这就保证了基础的低延迟。

在具体保障上,为了兼顾乱序数据和业务对时间的定义,Flink 提供了事件时间语义,就是说用数据本身发生的时间而不是机器接收到的时间来算窗口、做统计。为了让事件时间能正常工作,就引入了一个很重要的机制——Watermark,相当于告诉系统“比这个时间早的数据理论上不会再来了”。窗口算子就根据水位线来决定“差不多了,可以触发计算往下发了”,这样既不会无限制等下去,也不会因为网络抖动丢三落四。

到了 SQL API 这边,用户不用自己写水印生成逻辑,直接在建表 DDL 里通过 WATERMARK FOR 声明一下某个时间字段,再指定一下最大乱序容忍度就行了,引擎会自动帮你生成水位线、驱动窗口触发,对使用者来说开箱即用,时效性不用你操心。

再补充一点,Flink 的反压机制和异步 I/O 也间接保证了时效性。反压让整个流不会被打崩,数据自己能“慢下来”但不会丢;异步访问外部系统时不会阻塞主处理流程,吞吐和延迟都有好处,两条 API 都能享受到。

然后是准确性,这个其实可以拆成结果正确和不丢不重。

Flink 通过状态和检查点(Checkpoint)来保证故障恢复后的计算结果正确,支持精确一次(exactly-once) 的状态一致性。DataStream API 里,你把中间结果放在 state 里,Flink 会定期异步做一个分布式快照,用 Chandy-Lamport 算法思想,失败后从最近的检查点恢复,中间没处理完的那部分数据配合可重放的 source 会重新消费,保证一个不丢、一个不重。

SQL API 更不用说,所有的聚合、join 的状态管理都是内置的,底层自动启用精确一次语义,用户连 checkpoint 的配置都不用写逻辑,只要开了 checkpoint,就能保证跑批和跑流结果一致。

端到端的精确一次还要看 sink 配合。Flink 提供了两阶段提交的能力,比如把结果写到 Kafka,可以用事务把 checkpoint 和提交绑在一起,要么全部生效,要么全部回滚,这样外部系统看到的也是准确一次的数据。这在 SQL 那边也能通过 connector 的配置做到,比如 'sink.semantic' = 'exactly-once' 之类的。

一句话总结: 时效性靠事件时间+水印+流式处理模型,SQL 帮你把水印细节封装好了,DataStream 让你更灵活定义。 准确性靠状态后端+检查点+精确一次语义,SQL 替你管好了所有状态,DataStream 让你可以精细控制状态逻辑,但最终可靠性保障源自同一个引擎。 所以在我眼里,这俩 API 只是表达方式不同,底层在时效性和准确性的保障上其实是同一套东西,不会因为选 SQL 就降级,也不会因为用 DataStream 就自动变精确,主要还是看怎么正确使用。

  • sql 语句通过java cc解析成AST(语法树),用 SqlNode 表示
  • 结合数字字典 catlog 去验证 sql 的语法
  • 将语法树转换成逻辑计划,用 relNode 表示
  • 然后优化逻辑计划,先基于 calcite rules 优化,再基于 flink定制优化rules 去优化
  • 将逻辑计划转成成Flink的物理执行计划
  • 最后调用相应的 tanslateToPlan 方法转换和利用 CodeGen 元编程成Flink的各种算子。
  1. 资源消耗:
  2. 内存方面:Flink CDC 需要在内存中维护一定的状态来跟踪维度表的变更,如记录已处理的事务位置等。若维度表数据量庞大或变更频繁,可能导致内存占用增加,若内存不足,可能引发频繁的垃圾回收,影响性能。
  3. CPU 方面:捕获和处理维度表变更数据需要一定的 CPU 资源。解析数据库日志、转换数据格式等操作都要消耗 CPU,可能导致 CPU 使用率上升,若系统 CPU 资源紧张,会影响 Flink 作业及其他相关系统的性能。
  4. 数据延迟:
  5. Flink CDC 捕获维度表变更数据存在一定延迟。从数据库发生变更到 Flink 作业获取到变更数据,有数据传输、处理等环节,若网络不稳定或处理流程复杂,延迟会增加,导致维度表数据不能及时更新,影响实时计算结果的准确性。
  6. 并发处理能力:
  7. 若多个 Flink 作业同时通过 Flink CDC 监控同一维度表,可能对数据库产生较大的并发读取压力。数据库的并发处理能力有限,可能导致查询性能下降,甚至出现阻塞,影响维度表的正常使用和 Flink 作业对维度数据的获取。
  • windowAll()函数,用于没有keyBy过的流,数据发送给下游的单个实例(下游并行度为1)
  • window()函数,用于经过keyBy的流,将数据流中的元素分配到相应的窗口中
  • Flink为我们提供了一些内置的WindowAssigner,分别对应滚动窗口、滑动窗口、会话窗口
  • trigger()函数,指定触发器Trigger(可选)
  • evictor()函数,指定清除器(可选)
  • .reduce/aggregate/process() 窗口处理函数,这里定义数据处理逻辑
  • 方式1:使用批处理执行环境(推荐)
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();// 创建批处理执行环境
DataSet<String> text = env.readTextFile("hdfs://path/to/file");// 读取数据源(有界数据)
DataSet<Tuple2<String, Integer>> counts = text
    .flatMap(new Tokenizer())
    .groupBy(0)// 批处理操作
    .sum(1);
counts.writeAsText("hdfs://path/to/output");// 输出结果
env.execute("Batch WordCount Example");// 执行作业
  • 方式2:使用流处理环境处理有界数据
// 创建流执行环境但设置执行模式为BATCH
StreamExecutionEnvironment env = StreamExecutionEnvironment
    .getExecutionEnvironment()
    .setRuntimeMode(RuntimeExecutionMode.BATCH); // 关键设置
DataStream<String> text = env.readTextFile("hdfs://path/to/file"); // 读取有界数据源
DataStream<Tuple2<String, Integer>> counts = text
    .flatMap(new Tokenizer())
    .keyBy(value -> value.f0)
    .sum(1); // 使用流式API处理
counts.writeAsText("hdfs://path/to/output"); // 输出结果
env.execute("Streaming API for Batch"); // 执行作业
  • 方式3:通过Table API/SQL实现
// 创建表环境(设置为批模式)
EnvironmentSettings settings = EnvironmentSettings
    .newInstance()
    .inBatchMode()  // 批处理模式
    .build();
TableEnvironment tableEnv = TableEnvironment.create(settings);
tableEnv.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)"); // 注册表
Table result = tableEnv.sqlQuery(
    "SELECT user, SUM(amount) FROM Orders GROUP BY user");// 执行批查询
result.executeInsert("ResultTable");// 输出结果

Flink的TopN的代码实现

  1. 基于 ProcessAllWindowFunction 的实现(简单但非最佳性能)
// 定义 TopN 算法
主函数:
    // 1. 创建流执行环境
    env = 创建流执行环境()

    // 2. 模拟数据源:包含(类别, 值)的元组
    数据源 = [
        (A, 10.0),
        (A, 30.0),
        (B, 5.0),
        (A, 20.0),
        (B, 15.0),
        (A, 40.0),
        (B, 25.0)
    ]

    // 3. 构建数据处理流水线
    数据源
        .创建全窗口滑动窗口(窗口大小=10秒, 滑动间隔=5秒)
        .处理(TopN处理函数(topN=2))
        .打印输出()

    // 4. 执行作业
    执行作业("TopN 示例")


// 定义 TopN 处理函数类
类 TopN处理函数:
    私有属性 topN

    // 构造函数
    构造函数(topN):
        本对象.topN = topN

    // 核心处理方法
    方法 处理(窗口上下文, 元素集合, 输出收集器):
        // 1. 初始化最小堆(按值升序)
        最小堆 = 新建优先队列(比较器=按元组第二个值比较)

        // 2. 遍历所有元素,维护堆的大小不超过 topN
        对于 元素 属于 元素集合:
            最小堆.加入(元素)
            如果 最小堆.大小 > topN:
                最小堆.取出最小元素()

        // 3. 按倒序收集结果(因为是最小堆,依次取出堆顶就是从大到小)
        结果数量 = 取最小值(最小堆.大小, topN)
        结果数组 = 新建字符串数组(结果数量)
        从 结果数量-1 到 0 倒序遍历 i:
            元素 = 最小堆.取出最小元素()
            结果数组[i] = 元素.类别 + ":" + 元素.值

        // 4. 输出最终结果
        输出收集器.收集("窗口 [" + 窗口上下文.窗口 + "] 的Top" + topN + ": " + 结果数组.用逗号连接)
  1. 基于 KeyedProcessFunction 的更优实现(性能更好)
// 定义 TopN 算法(按 key 分区后实时维护 TopN)
主函数:
    // ... 省略创建env和数据源(与示例1相同)

    // 按类别分区,然后使用 KeyedProcessFunction 处理
    数据流
        .keyBy(元素 -> 元素.类别)  // 按类别分区
        .process(TopN键控处理函数(topN=2))
        .打印输出()

    // 执行作业
    执行作业("TopN Keyed 示例")

// 定义 TopN 键控处理函数类
类 TopN键控处理函数 继承 KeyedProcessFunction:
    私有属性 topN
    私有属性 状态  // 每个key维护一个最小堆状态

    // 构造函数
    构造函数(topN):
        本对象.topN = topN

    // 打开函数:初始化状态
    方法 打开(配置参数):
        // 为每个key创建一个ValueState,存储优先队列(最小堆)
        状态描述符 = 新建ValueState描述符("topN-queue", 优先队列类型)
        状态 = 运行时上下文获取状态(状态描述符)

    // 处理每个元素
    方法 处理元素(当前元素, 上下文, 输出收集器):
        // 从状态中获取当前key的最小堆
        堆 = 状态.值()
        如果 堆 为空:
            堆 = 新建优先队列(比较器=按元素的值升序排序)

        // 将新元素加入堆
        堆.加入(当前元素)

        // 如果堆大小超过 topN,移除最小元素
        当 堆.大小 > topN:
            堆.移除并返回最小元素()

        // 更新状态
        状态.更新(堆)

        // 如果需要定时输出,可以在这里注册定时器
        // 否则可以直接在这里整理结果并输出

    // 定时器触发时输出结果(可选)
    方法 定时器触发(时间戳, 上下文, 输出收集器):
        // 从状态取出堆
        堆 = 状态.值()
        如果 堆 为空:
            返回

        // 按从大到小收集结果(因为是最小堆,倒序取出就是降序)
        结果数组 = 新建列表
        堆大小 = 堆.大小
        对于 i 从 0 到 堆大小-1:
            结果数组.插入(0, 堆.取出最小元素())  // 每次插在头部,自然就是降序

        // 输出结果
        输出收集器.收集("Key " + 当前key + " 的Top" + topN + ": " + 结果数组.用逗号连接)

        // 清空状态(如果需要)
        状态.清除()
  1. 基于 Window 和 Reduce 的高效实现(推荐)
// 定义 TopN 算法(预聚合优化,性能最佳)
主函数:
    // ... 省略创建env和数据源(与示例1相同)

    // 按类别分区,开滑动窗口,先预聚合每个key的最大值,再输出TopN
    数据流
        .keyBy(元素 -> 元素.类别)      // 按类别分区
        .时间窗口(窗口大小=10秒, 滑动间隔=5秒)
        .reduce(最大值预聚合器(), TopN窗口处理函数(topN=2))
        .打印输出()

    // 执行作业
    执行作业("TopN Window 示例")


// 最大值预聚合器(每个key在窗口内只保留最大值,减少数据量)
类 最大值预聚合器 实现 ReduceFunction:
    方法 聚合(value1, value2):
        // 返回值较大的那个元素
        如果 value1.值 > value2.值:
            返回 value1
        否则:
            返回 value2


// 计算TopN的窗口函数
类 TopN窗口处理函数 
    继承 ProcessWindowFunction 
    实现 CheckpointedFunction:

    私有属性 topN
    私有属性 顶层项目状态  // 状态存储,用于容错

    // 构造函数
    构造函数(topN):
        本对象.topN = topN

    // 窗口处理方法
    方法 处理(key, 上下文, 预聚合后的最大值集合, 输出收集器):
        // 1. 收集所有类别对应的最大值(经过reduce预聚合后,每个key只有一个最大值)
        所有顶层项目 = 新建列表
        对于 项目 属于 最大值集合:
            所有顶层项目.添加(项目)

        // 2. 按值降序排序
        所有顶层项目.排序(比较器=按值降序)

        // 3. 取前 topN 个
        数量 = 取最小值(所有顶层项目.大小, topN)

        // 4. 整理并输出结果
        结果列表 = 新建列表
        对于 i 从 0 到 数量-1:
            项目 = 所有顶层项目[i]
            结果列表.添加(项目.类别 + ":" + 项目.值)

        输出收集器.收集("窗口 [" + 上下文.窗口 + "] 的Top" + topN + ": " + 结果列表.用逗号连接)

    // 保存状态快照(Checkpoint)
    方法 保存快照状态(上下文):
        // 完整实现需要将状态持久化,这里省略

    // 初始化状态
    方法 初始化状态(上下文):
        // 完整实现需要从持久化状态恢复,这里省略
  • 选择哪种实现?
  • 小数据量 & 简单实现:使用第一种基于 ProcessAllWindowFunction 的方法
  • 实时TopN & 高效状态管理:使用第二种基于 KeyedProcessFunction 的方法
  • 窗口TopN & 最佳性能:使用第三种基于 Window + Reduce + ProcessWindowFunction 的组合方法
  • 优化技巧
  • 预聚合:在处理大量数据时,先在本地进行预聚合
  • 状态管理:使用 Flink 的状态管理机制确保容错性和一致性
  • 高效数据结构:使用最小堆而不是全排序
  • 键分区策略:合理设计 keyBy 策略避免数据倾斜
  • Watermark处理:正确处理事件时间水印,处理延迟数据

项目

迟到处理: Flink业务默认基于event time时间戳,处理迟到数据时,会根据watermark进行处理。 正常数据则正常进入窗口聚合 晚于窗口长度的数据,但是在窗口关闭延迟内,会进入窗口重新计算 超过窗口关闭延迟的数据则默认是直接丢弃,但是在实际环境中可以配置输出到侧输出流

乱序处理: Checkpoint + 状态 TTL:流计算内部乱序、迟到兜底

核心机制可分为全量快照阶段和增量日志阶段,并通过 Flink 检查点(Checkpoint) 保证精确一次(Exactly-Once)语义。

  1. 无锁全量快照(Snapshot)—— StartupOptions.initial() 的首次执行 传统全量同步会用 SELECT * 加全局读锁(FLUSH TABLES WITH READ LOCK),影响业务。Flink CDC 采用无锁分块(Chunking)算法: 分块读取:根据主键(或唯一索引)将表数据切分成多个数据块(Chunk),并行读取。 一致性点(Consistency Point):开始快照前,先获取当前 MySQL 的 binlog 位点(GTID 或 File+Position),记录快照起始的日志偏移量。 可重复读(RR)隔离:利用 MySQL 的 REPEATABLE READ 事务隔离级别,保证每个分块读取到的数据是事务开始时的快照,不受其他事务更新影响,从而避免全局锁表。 断点续传:每个分块读取完成后,Flink 会记录该分块的完成状态,即使任务重启,未完成的分块可重新读取。

  2. 增量日志阶段(Binlog Streaming)—— 快照完成后的实时监听 全量快照完成后,Flink CDC 会转为 binlog 流式消费,原理等同于 MySQL 的 SHOW MASTER STATUS 和 COMREGISTERSLAVE / START SLAVE: 伪装成 MySQL Slave:Flink CDC 连接器伪装成 MySQL 的从库,向 MySQL 主库发送 DUMP BINLOG 请求。 解析 Binlog 事件:Debezium 引擎解析 binlog 中的 WRITEROWS、UPDATEROWS、DELETE_ROWS 事件,并根据表结构将其转换为 Flink 的 SourceRecord 数据。 无缝衔接:增量模式从全量快照开始时记录的 binlog 偏移位点 开始消费,保证全量与增量之间无数据丢失、无重复(即快照记录是位点之后的变更)。

  3. 状态持久化与断点续传(Exactly-Once 保证) Flink CDC 在 Flink 的 Source 层面实现了分区(Split)状态管理: 全量快照中,每个 Chunk 的读取进度(当前主键偏移量)和 binlog 当前位点 都会作为状态存入 Flink 的 状态后端(如 HDFS)。 配合 Flink 的 Checkpoint 机制,定期将状态持久化。

任务故障重启时: 未完成的全量 Chunk 会从 checkpoint 记录的偏移量继续读取; 若已进入增量阶段,则从 checkpoint 保存的 binlog 位点继续消费(通过 SET @masterbinlogchecksum 等机制定位),保证数据既不丢失也不重复。

  1. 数据结构转换(JsonDebeziumDeserializationSchema) Debezium 原始输出是复杂的 SourceRecord 结构(包含 before、after、op、source 元数据)。项目中指定的反序列化器将其转换为精简的 JSON 字符串

首先,我的这个实时数仓项目里,状态说白了就是在算子里面维护一个“记忆”,因为流数据是源源不断的,得记住前面来过什么数据,后面的计算才有依据。我们项目里主要用了两种:Keyed State(键控状态) 和广播状态。

我挑几个典型的场景说明一下:

第一,修复日志里的新老访客标记(数据清洗)。 前端埋点上报的 is_new 字段经常不准,比如老用户清了缓存就成新用户了。我们当时就用了一个 ValueState,按照设备ID(mid)分组,存了每个设备的“首次访问日期”。 逻辑很简单:如果日志说他是新用户,但我状态里存的首次访问日期不是今天,那我立马把他纠正成老用户。这样就靠状态把前端脏数据给“拨乱反正”了。

第二,做 UV(独立访客)和会话数统计。 这也是靠状态去重。比如算 UV,我还是用 ValueState 存这个用户最后一次访问的日期,每次来一条数据就比对一下,只要状态里的日期不是今天,UV 才加 1,然后更新状态。这样就保证了同一天内同一个用户只算一次。 会话数(SV)的时候,我怕状态无限膨胀,特意给状态设置了 1 小时的 TTL(存活时间),会话结束了状态自动就清掉了,防止内存爆掉。

第三,处理订单数据的回撤流(去重)。 这个是我们比较巧妙的一个点。DWD 层订单明细因为有 left join,会产生回撤数据。我们没傻等窗口结束,而是用状态存了上一条数据。当新数据来了,我们先把状态里上一条数据的金额取反,发下去做冲销,再把新数据更新到状态里。这就跟做账似的,不用等窗口关,实时就能输出精准结果,时效性特别好。

第四,也是我们项目的一个亮点——用广播状态做动态配置。 我们的维度表和事实表分流规则是写在 MySQL 配置表里的。我们通过 Flink CDC 监听这个配置表,把它做成一个广播流,跟主流数据做 connect 连接。 这样一来,配置一但改了,广播状态里的数据就变了,下游任务立马就能感知到。不像以前改个配置还得重启任务,现在完全热加载,特别灵活。

最后,我再补充一个点,虽然不算 Flink 自带状态,但我们配合了 Redis 旁路缓存。我们查 HBase 维度数据的时候,先查 Redis,如果命中就直接用。同时我们有个联动机制:当 DIM 层发现 HBase 里的维度数据被更新或删除了,我们会在 Sink 函数里手动删掉对应的 Redis 缓存。这样就保证了缓存数据和源头永远一致。

状态的生命周期?

状态生命周期的重要性(防止 OOM、处理延迟数据)

第一,也是最常用的,就是利用 Flink 自带的 StateTtlConfig 设置过期时间。 我们主要分了这么几种情况: 短生命周期的状态(比如 60 秒):在交易域下单统计(DwsTradeSkuOrderWindow)里去重的时候,因为订单的回撤流通常几秒内就来了,所以我们给状态设置了 60 秒的 TTL。这样数据一旦处理完,状态立马就释放了,不会长时间占着内存。 中等生命周期的状态(比如 1 小时):在统计会话数(SV)的时候,我们给状态设置了 1 小时的 TTL。因为正常用户的会话不可能持续一整天,1 小时不活动基本就断了,状态自动清理,防止内存里堆满无效的 session 信息。

第二,在 Flink SQL 的流表关联(双流 Join)中,我们设置了空闲状态保留时间。 在做订单明细宽表的时候,我们用 SQL 做了多表 Join。如果不设置,Flink 会默认永久保留左右表的状态,这在大数据场景下很容易 OOM。我们在 BaseSQLApp 里通过 tEnv.getConfig().setIdleStateRetention(Duration.ofSeconds(5)) 指定了空闲状态最多存活 5 秒。这意味着如果 5 秒内没有更新的订单数据到来,这条订单关联的状态就会被淘汰,既保证了数据能对上,又释放了压力。

第三,对于广播状态(配置流),我们不设置 TTL。 像 DIM 层的维度配置和 DWD 层的分流配置,我们是靠 CDC 读取 MySQL 配置表写进广播状态的。这种状态我们不给它设置 TTL,默认永久保留。因为配置必须一直在那里,除非我们在配置表里执行了 Delete 操作,代码里会显式地调用 state.remove(key) 手动删除,保证了配置的强一致性。

第四,对于“新老访客修复”这种业务强相关的状态,我们依赖业务逻辑更新,不依赖 TTL。 比如判断用户是否今日首次访问,状态里存的是“日期”。我们不需要设置 TTL 去删除它,因为当日期变成明天时,业务逻辑自己就会判断出“不等于今天”,然后自动覆盖旧状态。这种状态的生命周期相当于“一天”,是被业务时间驱动更新的。 最后,除了数据的 TTL,我们也特别关注检查点(Checkpoint)的超时机制。 比如我们在 BaseApp 里设置了 setCheckpointTimeout(10000)(10秒超时)。如果这次持久化状态太慢或者失败了,我们宁可让它失败重启,也不要让旧状态卡在那里导致反压。这也算是从故障恢复的视角控制了状态的“有效性窗口”。

总的来说,我们管理状态生命周期的核心原则就是:绝不保留无用数据,能短则短,按需设置,配合业务逻辑来驱动状态的更新和淘汰。 这样既保证了 Exactly-Once 的准确性,也保障了集群的稳定性。

// 统计每分钟日志数量(滚动窗口)
主函数:
    // 1. 创建流执行环境
    env = 创建流执行环境()

    // 2. 设置水印(事件时间需要)
    env.设置水印允许延迟(允许延迟时间=5秒)

    // 3. 读取数据源(例如 Kafka)
    数据流 = env.addSource(Kafka数据源)
        .map(记录 -> 解析日志, 提取时间戳)

    // 4. 按分钟开窗聚合,统计每一分钟的日志数量
    结果流 = 数据流
        .keyBy(日志 -> 日志.级别)  // 可选:按日志级别分组统计
        .滚动窗口(时间大小=1分钟)  // 滚动窗口,每分钟一个窗口
        .aggregate(计数聚合器())     // 聚合计算数量

    // 5. 输出结果(例如写到 Elasticsearch)
    结果流.addSink(结果输出)

    // 6. 启动作业
    env.execute("每分钟日志数量统计")


// 计数聚合器
类 计数聚合器 实现 AggregateFunction:
    // 创建累加器,初始计数为0
    方法 创建累加器():
        返回 0

    // 每来一条元素,计数加1
    方法 添加(元素, 累加器):
        返回 累加器 + 1

    // 返回最终结果
    方法 获取结果(累加器):
        返回 累加器

    // 合并两个分区的累加器(会话窗口等需要)
    方法 合并(累加器1, 累加器2):
        返回 累加器1 + 累加器2

核心步骤:读取数据源 → 提取时间戳/生成水印 → keyBy分区 → 定义窗口 → 聚合计算 → 输出结果

Flink CDC 在本项目中扮演了配置中心动态感知的角色,通过监听 MySQL 配置表的变更,实现了维度表和事实表分流的配置热更新,大幅提升了实时数仓的灵活性和可维护性。

(1)DIM 层动态维度分流(DimApp) 目的:从 MySQL 配置表 tableprocessdim 读取维度表的分流规则(如源表、目标 HBase 表、列族、RowKey、输出字段等),并实时感知配置变更。

实现方式: 使用 MySqlSource 构建 CDC 源,监听 gmall2023config.tableprocessdim 表。 反序列化使用 JsonDebeziumDeserializationSchema,将变更事件转为 JSON 字符串。 启动模式设置为 StartupOptions.initial(),表示首次启动时全量读取配置,之后通过 binlog 实时捕获增删改。 将配置流作为广播流,与主流(Kafka 的 topicdb 业务数据)进行 connect,通过 BroadcastProcessFunction 将最新配置存入广播状态,同时预加载配置到内存(防止配置流晚于数据流到达导致数据丢失)。 主流数据根据广播状态中的配置信息,决定是否保留、删除多余字段,并写出到 HBase 对应的维度表。

(2)DWD 层事实表动态分流(DwdBaseDb) 目的:从配置表 tableprocessdwd 读取事实表的分流规则(源表、操作类型、目标 Kafka 主题、输出字段),实现多条事实表的动态分流(如优惠券领取、收藏、注册等)。 实现方式:与 DIM 层类似,也是通过 FlinkCDC 读取配置表,将配置流作为广播流,与 topic_db 主流关联,根据配置动态过滤字段并写入不同 Kafka 主题。 关键逻辑与 DIM 层完全一致,只是配置表实体类和目标存储(Kafka)不同。

  1. 配置数据源:在 Flink 作业中配置要监控的维度表所在的数据库连接信息,包括数据库类型、主机地址、端口、用户名、密码等,以及指定要监控的具体维度表。
  2. 启用 CDC:使用 Flink CDC 相关的 API 或连接器,启用对维度表的 CDC 功能,开始捕获表中的数据变更。
  3. 处理变更数据:在 Flink 作业中,可以定义对捕获到的变更数据的处理逻辑。例如,当维度表中的某条记录发生更新时,可以将新的维度信息与其他业务数据进行关联,或者根据变更情况更新缓存中的维度数据,以便在后续的计算中使用最新的维度信息。
  4. 参考回答:Flink CDC基于数据库日志(如MySQL的binlog)捕获变更事件,通过Debezium解析为流数据。为保证一致性,我采用以下策略:
    • 使用HBase的原子操作(如CheckAndPut)实现幂等写入;
    • 结合Flink的Exactly-Once语义和两阶段提交(2PC)确保端到端一致性。

CDC(Change Data Capture,变更数据捕获)的底层原理主要基于数据库事务日志的解析,其核心思想是实时捕获数据库中的变更操作(如INSERT/UPDATE/DELETE),并将这些变更以有序事件流的形式输出,供下游系统消费。以下是其技术原理的分步解释:


  1. 数据库事务日志(Transaction Log)
  2. 作用:所有数据库(如MySQL、PostgreSQL、Oracle等)都会维护事务日志(如MySQL的binlog、PostgreSQL的WAL),记录每一次数据变更的详细信息。

  1. 日志解析与事件生成
  2. CDC工具(如Debezium、Canal、Maxwell)会:
  3. 监听日志:持续监控数据库的日志写入位置(如MySQL的binlog offset)。
  4. 解析日志:将二进制日志转换为结构化事件,例如:
  5. 事件排序:基于事务时间戳(LSN, Log Sequence Number)保证事件顺序。

  1. 数据传输与消费
  2. 输出目标:
  3. 消息队列(如Kafka、RabbitMQ):供实时流处理引擎消费。
  4. 数据仓库/湖(如Snowflake、S3):用于离线分析。
  5. 数据库(如PostgreSQL、Elasticsearch):实现跨系统数据同步。
  6. 消费模式:
  7. At-Least-Once:确保事件不丢失(可能重复)。
  8. Exactly-Once:通过事务或幂等设计实现精确一次消费。

如果Flink作业出现反压(Backpressure),可能的排查步骤是什么?

首先,我不会一上来就盯着Source看,因为反压一定是下游先堵住的。 我第一步肯定是打开Flink的Web UI,直接看哪个算子的Subtask是红色的“HIGH”反压状态。关键技巧是——反压链条里,最下游那个标红的算子就是真正的瓶颈,它前面的所有反压都是被传导的。找到它之后,我会点进去看它的输入输出速率,如果输入远大于输出,那问题就锁死在这了。

定位到瓶颈算子后,第二步我会看它的资源消耗。 先看CPU,如果CPU已经跑满,那说明计算逻辑太重,可能需要优化代码或提高并行度;但如果CPU很低,那就得小心了——大概率是IO等待或者网络堵了。我会顺便瞄一眼GC时间,如果频繁Full GC,那得先调JVM内存,否则其他排查都没意义。

第三步,我最怕遇到数据倾斜,因为隐蔽性最强。 我会在UI里对比瓶颈算子各个并行子任务的numRecordsIn,如果某个子任务处理的数据量明显比别人多好几倍,那肯定就是KeyBy的key分布不均。这种时候我会考虑在KeyBy前加随机前缀打散,或者临时用Rebalance重分区先缓解。

如果资源没问题,倾斜也不明显,那就要看业务逻辑了。 比如如果是异步IO算子,我会怀疑外部存储(比如Redis或HBase)响应变慢了,得检查连接池和超时时间;如果是窗口算子,我会看窗口触发频率是不是太高,导致状态读写压力大;如果是自定义函数,得检查有没有加锁或同步调用外部接口——这些都很容易卡住。

接着,我会查状态后端,尤其是用RocksDB的时候。 我会关注RocksDB的读写延迟和Compaction指标,如果Compaction太频繁,会阻塞写入,这时候可以调大Compaction线程数或者开启增量检查点。同时也会看状态大小,如果状态膨胀到几个G,每次访问都走磁盘,那肯定慢,得优化状态结构。

最后,千万别忘了Sink端。 很多时候反压源头其实是写入目标系统慢了,比如Kafka分区负载不均、HBase RegionServer卡顿、或者MySQL写入限流。我会看Sink算子的重试次数和发送延迟,如果频繁重试,那就得在Sink端开启批量写入,调大flush间隔,或者跟DBA确认下外部系统状态。

如果上面这些都查了还没头绪,我会用终极手段——抓火焰图。 在TaskManager上跑一下async-profiler,看线程堆栈卡在哪个具体方法上,往往能直接定位到某行序列化代码或者网络读取。

当然,如果线上已经严重积压,我会先止血。 比如临时增大网络内存占比,或者调大Checkpoint间隔(因为Barrier对齐也会加剧反压),如果Flink版本支持,我甚至会直接动态调整并行度重启。

总之,我的排查思路就是从UI定位瓶颈,逐层往下剥,先硬件后业务,最后再考虑外部系统和代码细节。

ReplacingMergeTree和普通的MergeTree有什么区别

在MergeTree中,虽然有主键,但是它没有唯一键的约束,写入数据的主键是可以重复的 ReplacingMergeTree可以在合并分区时,删除重复的数据条目

评论