很多用SeaTunnel做数据同步的小伙伴,都会遇到一个头疼的问题:明明源端配置了好几路并发拉取数据,结果下游写入的时候要么报错,要么CPU/内存被占满,其实这大概率是源端的拉取速度和下游的写入速度不匹配,这时候SeaTunnel的背压机制就会悄悄生效,帮你解决这个问题。
一、先搞懂SeaTunnel里的“拉取-写入”是啥
1.1 简单说下SeaTunnel的干活逻辑
就像你要把A地的一批快递搬到B地,SeaTunnel就是那个搬家的车队:源端是A地的快递仓库,你可以配置“并发”就是同时开几辆货车去仓库打包,每辆车负责拉取一部分数据;下游是B地的驿站,就是你要把快递送到的地方,每辆车拉到的快递要交到驿站,驿站的处理速度就是下游的写入速率——比如驿站只有1个工作人员,每秒只能处理10个快递,结果仓库每秒钟发20个,那就会堆积,甚至快递丢了。
1.2 为啥会有速率不匹配的问题?
这种情况太常见了,比如源端是MySQL,配置了3个并发拉取,每秒钟能拿到1000条数据,但下游是Kafka或者Elasticsearch,因为索引优化或者节点性能问题,每秒钟只能写500条,或者下游是API接口,因为限流规则最多每秒接收300条,这时候就会出现“快的拉取、慢的写入”,数据越积越多,要么内存炸了,要么写入失败。
二、当快慢不搭时,背压是怎么起作用的?
背压这个词听起来专业,其实说白了就是“自适应限流”,就像你用杯子接自来水,水龙头开太大,杯子倒不赢,水会溢出来,背压就是让水龙头开小一点,或者杯子先存一点水,等倒完再继续,不让溢出。
2.1 举个接地气的SeaTunnel示例
我们用SeaTunnel的配置来演示这个过程,这个示例的技术栈就是SeaTunnel,全程只用这一个技术栈:
# SeaTunnel配置文件,演示背压机制的生效逻辑
env {
parallelism = 3 # 源端拉取的并发数,相当于3辆货车同时从仓库拉货
job.mode = "streaming" # 流处理模式,模拟持续不间断的拉取-写入过程
backpressure.enable = true # 核心配置:开启背压机制,默认是开启的但最好显式写出来
backpressure.max.block.size = 1024 # 缓存数据的最大体积(MB),相当于驿站临时堆货的最大空间,超了就暂停拉货
}
source {
MySQL-CDC {
plugin_output = "sync_data"
url = "jdbc:mysql://localhost:3306/test_order_db"
username = "sync_user"
password = "sync_pass"
table-names = ["order_info"]
parallelism = 3 # 和env的parallelism对应,源端并发拉取的线程数
}
}
sink {
Kafka {
plugin_input = "sync_data"
bootstrap.servers = "localhost:9092"
topic = "sync_order_topic"
# 下面的配置模拟下游写入慢,故意加延迟和批次限制,让速率不匹配
properties.acks = "1" # 只等Leader节点确认,比全副本确认快但还是有延迟
properties.batch.size = 16384 # 16KB的批次,批次越大,单次写入越慢
properties.linger.ms = 5000 # 延迟5秒再发送,相当于让下游慢慢处理,故意变慢
}
}
这个配置里,下游Kafka被我们调慢了,这时候SeaTunnel的背压机制就会启动:它会实时监控内部缓存的待写入数据量,当缓存超过1024MB时,就会给源端的拉取线程发暂停信号,让拉取线程不再拿新数据,等下游把缓存里的1024MB数据慢慢写入Kafka,缓存降到安全值以下,源端的拉取线程再恢复干活,这样就不会有数据堆积。
2.2 背压生效的细节
背压不是直接停掉所有拉取,而是“局部限流式暂停”:比如源端有3个并发拉取线程,其中2个已经把数据推到缓存,1个还在 MySQL里拉,这时候如果缓存快满了,正在拉的那个线程会先停,已经推到缓存的线程,会等下游处理完再推,相当于“让跑最快的那辆车先停一停,等慢的跟上”,整个过程是自动的,不需要人工干预,也不会丢数据。
三、背压机制的适用场景
3.1 常见的需要背压的情况
背压不是万能的,但在这些场景下特别好用:第一,源端是高并发数据库,下游是低吞吐的存储,比如MySQL同步到Elasticsearch,ES的写入速度远不如MySQL的拉取速度;第二,流处理CDC同步,下游是消息队列有自带限流,比如Kafka的分区数不够导致写入慢;第三,批量同步大数据量,比如用SeaTunnel同步10亿条数据到ClickHouse,没开背压会让内存被数据占满崩溃,背压能保护系统。
3.2 什么时候要小心背压?
也不是所有情况都要开背压,比如同步小数据量,源端和下游速度匹配,开背压反而会增加延迟,因为暂停的时间相当于额外的同步时间;还有实时性要求极高的场景,比如直播数据同步,延迟要求在1秒以内,背压的暂停会导致延迟升高,这时候应该调小源端的并发数,而不是依赖背压。
四、背压的优缺点和注意事项
4.1 优点
最核心的优点就是“保护系统”:不会因为数据堆积导致SeaTunnel的内存溢出,不会让下游的存储或者服务被压垮而崩溃;其次是“自动适配”,不用你手动去算源端和下游的速率,背压会根据实时情况自动调整,节省了人工调优的成本。
4.2 缺点
背压的缺点也很明显:会增加数据同步的延迟,因为背压的暂停就是为了让下游赶上来,这段等待的时间就是延迟;还有参数配置不好会适得其反,比如把max.block.size设太小,缓存稍微涨一点就暂停,频繁的暂停和恢复反而会降低整体同步的吞吐量,设太大的话,缓存堆积太多,内存还是会有风险。
4.3 注意事项
用背压的时候要注意这几点:第一,背压默认是开启的,但旧版本的SeaTunnel可能关闭了,一定要显式配置backpressure.enable = true;第二,max.block.size要根据下游的性能来调,比如下游ES每秒能写10GB,那max.block.size设成5-10GB就差不多,太小会频繁暂停,太大有内存风险;第三,不要和源端的限流同时开,比如MySQL自己做了每秒最多查500条的限流,SeaTunnel再开背压,两个限流冲突会影响同步效率;第四,背压只适合流处理或者持续同步的场景,一次性批量同步的话,不需要背压,直接控制源端并发就行。
五、总结
其实SeaTunnel的背压机制就是自带的“智能刹车”,当你配置多源并发拉取,下游写不动的时候,它会自动踩刹车,不让数据堆积,既保护了SeaTunnel本身和下游服务,又能保证数据同步顺利进行。平时用的时候,只要记得显式开启背压,调对缓存大小,就不用再担心速率不匹配的问题,遇到同步报错、内存溢出这些问题,先检查下背压是不是开了,参数对不对,大概率能解决大部分性能瓶颈。
评论
围绕“浅谈SeaTunnel的源端并发拉取与下游写入速率不匹配时背压机制是如何生效的”参与讨论