最近做了一个Flink实时计算项目,一开始性能很差,延迟很高,经常背压。经过一轮深入的性能优化,把吞吐量提升了5倍,延迟降到了毫秒级。记录一下整个优化过程,从资源配置到代码优化,从状态管理到checkpoint,一次说清楚。
先说说项目背景。我们需要对用户的实时行为数据做流式计算,包括实时统计、实时推荐、实时风控等场景。数据量很大,每秒大概有几十万条数据,要求延迟在秒级以内,最好能做到毫秒级。我们选择了Flink作为实时计算引擎,因为它的流处理能力很强,支持事件时间、状态管理、exactly-once语义等特性。
项目一开始,我们按照常规的方式写了Flink作业,部署到集群上运行。结果发现性能很差,吞吐量上不去,延迟很高,而且经常出现背压(backpressure),数据处理不过来,堆积在队列里。最严重的时候,延迟达到了几分钟,完全满足不了业务需求。
于是我们开始做性能优化。经过大概两周的努力,从资源配置、代码优化、状态管理、checkpoint优化等各个方面做了调整,最终把吞吐量提升了5倍,延迟降到了100毫秒以内,背压问题也彻底解决了。今天把整个优化过程记录下来,分享给大家。
一、性能分析:先搞清楚瓶颈在哪里
性能优化的第一步,不是上来就改代码,而是先做性能分析,搞清楚瓶颈到底在哪里。没有数据的优化就是瞎调,很可能改了半天发现改的地方根本不是瓶颈。
Flink提供了很完善的监控指标,可以通过Flink Web UI看到每个算子的吞吐量、延迟、背压情况、状态大小、checkpoint时间等。我们通过这些指标,发现了几个主要的瓶颈:
第一个瓶颈是数据源的读取速度。我们的数据源是Kafka,一开始我们设置的Kafka消费者并行度比较低,只有4个,而Kafka的分区有16个。这样导致Kafka的消费速度跟不上,数据在Kafka里堆积,Flink作业虽然有处理能力,但拿不到数据。
第二个瓶颈是某个关键算子的处理速度慢。我们的作业中有一个富函数(RichFlatMapFunction),里面做了比较复杂的逻辑处理,包括数据清洗、格式转换、维度关联等。这个算子的并行度设置得不够高,而且代码里有一些性能问题,导致这个算子的处理速度跟不上上游的数据输入速度,成为了整个作业的瓶颈,上游出现了背压。
第三个瓶颈是状态访问慢。我们的作业用到了键控状态(Keyed State),用来存储用户的实时统计信息。一开始我们用的是默认的状态后端(MemoryStateBackend),状态都存在JVM堆内存里。随着数据量的增加,状态越来越大,JVM堆内存不够用了,导致频繁的Full GC,作业卡顿严重,处理速度上不去。
第四个瓶颈是checkpoint时间太长。我们开启了checkpoint来保证exactly-once语义,一开始设置的checkpoint间隔是1分钟。但因为状态比较大,每次checkpoint都要花好几分钟,而且checkpoint的时候会暂停数据处理,导致延迟增加。有时候checkpoint还没完成,下一次checkpoint又开始了,形成了恶性循环。
找到了这几个瓶颈之后,我们就开始针对性地做优化了。
二、资源配置优化
第一个优化是资源配置,包括并行度设置、内存配置、CPU配置等。
首先是并行度设置。Flink作业的并行度决定了每个算子有多少个并行实例来处理数据。并行度设置得太低,处理能力不够,会出现瓶颈;设置得太高,会浪费资源,而且增加调度和数据交换的开销。
我们的优化原则是:每个算子的并行度要根据它的处理能力和数据量来设置,不能一概而论。对于数据源(Kafka消费者),并行度应该和Kafka的分区数一致,这样每个消费者实例对应一个分区,消费效率最高。我们把Kafka消费者的并行度从4改成了16,和Kafka分区数一致,消费速度立刻就上去了。
对于处理能力比较弱的算子(比如我们那个复杂的富函数),并行度要设置得高一些,让更多的实例来分担处理压力。我们把那个富函数的并行度从8改成了32,处理速度提升了好几倍。对于处理能力比较强的算子(比如简单的map、filter),并行度可以设置得低一些,不需要太高。
另外,Flink支持给不同的算子设置不同的并行度,不需要整个作业用同一个并行度。我们根据每个算子的实际情况,分别设置了合适的并行度,这样既能保证处理能力,又不会浪费资源。
其次是内存配置。Flink的内存配置比较复杂,分为JVM堆内存、堆外内存、托管内存等。不同的状态后端对内存的需求也不一样。
我们一开始用的是MemoryStateBackend,状态存在JVM堆内存里,所以给TaskManager分配了很多堆内存(8GB)。但即使这样,状态大了之后还是会频繁Full GC。后来我们改成了RocksDBStateBackend,状态存在磁盘上,堆内存只需要存一些缓存和运行时数据,就不需要那么多堆内存了。我们把堆内存降到了4GB,同时增加了堆外内存和托管内存的配置,给RocksDB用。调整之后,Full GC的问题解决了,作业运行稳定了很多。
另外,Flink的内存配置要根据实际情况来调整,不要盲目给太多内存。给太多堆内存反而会导致Full GC时间变长,影响性能。要根据状态后端、数据量、作业复杂度等因素,合理配置各个部分的内存。
第三是CPU配置。Flink作业对CPU的需求也比较大,尤其是计算密集型的作业。我们一开始给每个TaskManager分配了2个CPU核心,后来发现不够用,CPU使用率经常达到100%。我们把每个TaskManager的CPU核心数增加到了4个,同时调整了每个TaskManager上的slot数量,让每个slot有足够的CPU资源。调整之后,CPU使用率降到了60%左右,处理速度也提升了。
三、代码优化
第二个优化是代码优化。Flink作业的性能,很大程度上取决于代码写得好不好。代码写得不好,再怎么调资源也没用。
我们在代码优化中做了以下几个方面的改进:
第一个是避免在算子函数里做重复的初始化工作。比如,我们的富函数里,原来在open()方法里做了很多初始化工作,包括加载配置文件、初始化数据库连接、初始化缓存等。这些工作本来是对的,但我们发现,有些初始化工作其实可以在作业提交的时候就做好,不需要在每个并行实例里都做一遍。
比如,配置文件的加载,我们原来在每个算子的open()方法里都加载一次,其实完全可以在主函数里加载好,然后通过构造函数传给算子。这样就避免了每个并行实例都重复加载配置文件,节省了初始化时间,也减少了资源消耗。
再比如,一些不变的字典数据,我们原来在每个算子的open()方法里都从数据库加载一次,其实可以在主函数里加载好,然后广播给所有算子(Broadcast State),或者通过构造函数传进去。这样就避免了每个实例都去数据库加载数据,减少了数据库的压力,也加快了初始化速度。
第二个是减少序列化和反序列化的开销。Flink在数据交换的时候,需要把数据序列化成字节流,传输之后再反序列化。序列化和反序列化是有开销的,如果数据结构比较复杂,开销会很大。
我们的优化方法是:尽量使用Flink内置的序列化器(比如POJO、Tuple、基本类型),这些类型的序列化器效率很高。避免使用复杂的嵌套对象或者第三方库的对象作为数据流的类型,因为这些对象的序列化效率比较低。如果必须使用复杂对象,尽量实现一个高效的自定义序列化器。
另外,尽量减少在算子之间传输的数据量。比如,不需要的字段就不要往下游传,可以在map算子中把不需要的字段去掉。这样既减少了序列化的开销,也减少了网络传输的开销。
第三个是合理使用状态。Flink的状态管理是它的核心特性之一,但如果用得不好,会严重影响性能。
我们的优化方法是:
- 只存储必要的状态,不要把所有数据都存到状态里。比如,用户的实时统计信息只存最近一段时间的,历史数据可以存到外部存储(比如Redis、HBase)里,需要的时候再查。
- 合理设置状态的TTL(Time To Live),让过期的状态自动清理掉,避免状态无限增长。我们给用户的实时统计状态设置了24小时的TTL,超过24小时没有更新的状态会自动清理。
- 对于不需要精确语义的场景,可以使用增量聚合(AggregateFunction或者ReduceFunction),而不是把所有原始数据都存到状态里。增量聚合只存聚合结果,状态大小小很多,访问速度也快。
第四个是避免在算子函数里做耗时的操作。比如,在map或者flatMap函数里同步调用外部接口(比如数据库查询、HTTP请求),会导致这个算子的处理速度很慢,成为瓶颈。
我们的优化方法是:
- 对于需要频繁查询的维度数据,提前加载到内存缓存里,在算子函数里直接查缓存,而不是每次都去查数据库。
- 对于必须调用外部接口的场景,使用异步IO(Async I/O),让请求异步并发执行,而不是同步阻塞等待。Flink的Async I/O功能可以很好地解决这个问题,能大幅提升处理外部接口的吞吐量。
- 对于批量处理的场景,可以攒一批数据一起处理,而不是每条数据都处理一次。比如,攒100条数据一起写数据库,比每条数据都写一次数据库效率高很多。
第五个是合理使用窗口。Flink的窗口(Window)是流处理中常用的功能,但如果用得不好,也会影响性能。
我们的优化方法是:
- 尽量使用增量聚合函数(AggregateFunction或者ReduceFunction),而不是全量窗口函数(ProcessWindowFunction)。增量聚合函数是在数据到达的时候就开始聚合,窗口触发的时候只需要输出聚合结果,不需要缓存所有原始数据,状态小,处理快。全量窗口函数需要把窗口内的所有数据都缓存起来,窗口触发的时候再处理,状态大,处理慢。
- 合理设置窗口的大小和滑动步长。窗口太大,状态就大,处理慢;窗口太小,触发太频繁,开销大。要根据业务需求和数据量,合理设置。
- 对于事件时间的窗口,合理设置watermark的延迟。延迟设置得太大,窗口触发太晚,延迟高;设置得太小,会有很多迟到数据,需要处理迟到数据,开销大。要根据数据的实际乱序情况,合理设置。
四、状态后端优化
第三个优化是状态后端的选择和配置。Flink有三种内置的状态后端:MemoryStateBackend、FsStateBackend、RocksDBStateBackend。选择合适的状态后端,对性能影响很大。
MemoryStateBackend把状态存在JVM堆内存里,访问速度快,但状态大小受限于堆内存,而且不适合大状态的场景,因为会导致频繁Full GC。适合状态比较小、对延迟要求很高的场景。
FsStateBackend把状态存在文件系统(比如HDFS、S3)上,平时的状态访问还是在堆内存里,checkpoint的时候才写到文件系统。比MemoryStateBackend更可靠,但状态大小还是受限于堆内存。
RocksDBStateBackend把状态存在RocksDB里(本地磁盘),checkpoint的时候写到文件系统。状态大小不受堆内存限制,可以存很大的状态。而且RocksDB是一个嵌入式的键值存储,读写性能不错。适合状态比较大的场景。
我们一开始用的是MemoryStateBackend,因为状态小的时候速度快。但随着数据量增加,状态越来越大,堆内存不够用了,频繁Full GC。后来我们改成了RocksDBStateBackend,状态存在磁盘上,堆内存只存运行时数据,Full GC的问题就解决了。
改成RocksDBStateBackend之后,我们还做了一些配置优化:
- 增加了RocksDB的内存配置,包括写缓冲区、块缓存、布隆过滤器等。合理的内存配置能大幅提升RocksDB的读写性能。
- 开启了RocksDB的增量checkpoint。增量checkpoint只上传和上一次checkpoint相比变化的部分,不需要每次都上传全量状态,能大幅减少checkpoint的时间和网络开销。
- 合理设置RocksDB的压缩方式。RocksDB默认用的是Snappy压缩,压缩速度快,但压缩率一般。如果状态比较大,可以考虑用LZ4或者ZSTD压缩,压缩率更高,能减少磁盘空间和checkpoint上传的数据量。
- 把RocksDB的存储目录配置在本地的SSD磁盘上,而不是机械磁盘或者网络磁盘。SSD的随机读写性能比机械磁盘好很多,RocksDB是随机读写比较多的存储,用SSD能大幅提升性能。
五、Checkpoint优化
第四个优化是checkpoint的优化。checkpoint是Flink保证exactly-once语义的核心机制,但如果配置不好,会严重影响性能。
我们的checkpoint优化包括以下几个方面:
第一个是合理设置checkpoint间隔。checkpoint间隔太短,checkpoint太频繁,开销大,影响吞吐量;间隔太长,故障恢复的时候需要重放的数据多,恢复时间长。要根据业务需求和状态大小,合理设置。我们的状态比较大,把checkpoint间隔从1分钟改成了5分钟,checkpoint的开销大大减少,吞吐量提升了。
第二个是合理设置checkpoint超时时间。如果状态比较大,checkpoint可能需要比较长的时间,如果超时时间设置得太短,checkpoint会频繁超时失败。我们把checkpoint超时时间从10分钟改成了30分钟,确保checkpoint有足够的时间完成。
第三个是开启非对齐checkpoint(Unaligned Checkpoint)。Flink 1.11之后引入了非对齐checkpoint,它可以在有背压的情况下也能快速完成checkpoint,不需要等所有在途数据都处理完。非对齐checkpoint会把在途数据也保存到checkpoint里,checkpoint的数据量会大一些,但checkpoint时间会短很多。我们开启了非对齐checkpoint之后,在有背压的情况下,checkpoint也能正常完成,不会因为背压而超时。
第四个是合理设置checkpoint的并发数。默认情况下,Flink同时只允许一个checkpoint在运行。如果checkpoint间隔比较短,而checkpoint时间比较长,就会出现上一个checkpoint还没完成,下一个checkpoint又开始了的情况。可以适当增加checkpoint的并发数,允许多个checkpoint同时运行。但要注意,并发checkpoint会增加资源消耗,要根据实际情况设置。
第五个是开启checkpoint的外部化存储。默认情况下,Flink的checkpoint在作业取消的时候会被删除。如果开启了外部化存储,checkpoint会保留在外部存储中,即使作业取消了也不会删除,这样可以从checkpoint恢复作业。这个功能对于生产环境来说很重要,能提高作业的可靠性。
六、其他优化
除了上面说的几个主要优化点,我们还做了一些其他的优化:
第一个是任务链(Operator Chain)的优化。Flink会把可以链接在一起的算子链接成一个任务,这样可以减少线程切换和数据交换的开销。默认情况下,Flink会自动做任务链优化,但有时候因为算子的并行度不同或者其他原因,无法自动链接。可以通过startNewChain()或者disableChaining()来手动控制任务链,把应该链接在一起的算子链接起来,提高性能。
第二个是数据本地性优化。Flink在调度任务的时候,会尽量把任务调度到数据所在的节点,减少数据传输的开销。对于数据源是HDFS或者其他分布式存储的情况,数据本地性尤其重要。可以通过调整slot分配策略和任务调度策略,提高数据本地性。
第三个是反压(Backpressure)处理。反压是Flink作业中常见的问题,当某个算子的处理速度跟不上上游的输入速度时,就会出现反压,上游的算子会被阻塞,吞吐量下降。解决反压的根本方法是找到瓶颈算子,提升它的处理能力(增加并行度、优化代码等)。另外,可以通过调整缓冲区的大小和超时时间,来缓解反压的影响。
第四个是监控和告警。性能优化不是一次性的工作,需要持续监控作业的运行状态,及时发现和解决问题。我们搭建了完善的监控体系,通过Prometheus采集Flink的 metrics,用Grafana做可视化展示,配置了关键指标的告警(比如延迟、吞吐量、背压、checkpoint失败等)。这样,作业出现问题的时候,我们能第一时间发现,及时处理。
七、优化效果
所有优化做完之后,我们重新测试了作业的性能,结果如下:
- 吞吐量:从原来的每秒10万条提升到了每秒50万条,提升了5倍
- 延迟:从原来的几分钟降到了100毫秒以内,满足了业务的毫秒级延迟需求
- 背压:彻底解决了,所有算子的处理速度都能跟上,没有出现背压
- 稳定性:作业运行稳定,Full GC消失了,checkpoint正常完成,没有出现故障
- 资源使用:CPU和内存使用率都在合理范围内,没有出现资源瓶颈
这个结果超出了我们的预期,也证明了性能优化的重要性。同样的硬件资源,通过合理的配置和代码优化,性能可以提升好几倍。
八、一些经验和建议
最后,分享一些Flink性能优化的经验和建议:
第一,先测量再优化。不要上来就改代码,先通过监控指标搞清楚瓶颈在哪里,然后针对性地优化。没有数据的优化就是瞎调,很可能做无用功。
第二,从简单到复杂。先做简单的优化(比如调整并行度、内存配置),再做复杂的优化(比如代码重构、状态后端切换)。简单的优化往往能取得不错的效果,而且风险小。
第三,关注数据倾斜。数据倾斜是流处理中常见的问题,某个key的数据量特别大,导致某个并行实例的处理压力特别大,成为瓶颈。如果发现某个算子的某个并行实例特别忙,其他实例很闲,那很可能是数据倾斜了。解决方法包括:对key做加盐处理、两阶段聚合、使用局部聚合等。
第四,合理使用Flink的高级特性。Flink有很多高级特性,比如Async I/O、Broadcast State、Incremental Aggregation等,合理使用这些特性能大幅提升性能。但也要注意,不要为了用而用,要根据实际场景选择合适的特性。
第五,版本要新。Flink的版本更新很快,每个版本都会修复很多bug,增加很多新功能,优化性能。尽量用比较新的稳定版本,能享受到版本更新带来的性能提升和bug修复。当然,升级版本之前要做好测试,确保兼容性。
第六,文档和社区。Flink的官方文档很完善,遇到问题的时候先查文档。Flink的社区也很活跃,邮件列表、Stack Overflow、GitHub issue上都有很多讨论,遇到问题可以去搜一搜,大概率能找到解决方案。
九、写在最后
Flink是一个非常强大的实时计算引擎,但它的性能优化也是一个比较复杂的话题,涉及资源配置、代码优化、状态管理、checkpoint等方方面面。这次优化经历让我对Flink的原理和性能调优有了更深的理解,也积累了很多实战经验。
性能优化是一个持续的过程,不是一次优化完就完事了。随着业务的发展,数据量会增加,需求会变化,性能瓶颈也会变化。需要持续监控,持续优化,才能保证作业的稳定高效运行。
希望我的这些优化经验能对大家有帮助。如果有什么问题或者不同的看法,欢迎在评论区交流。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录