news 2026/8/10 14:57:13

当 Redis 写入成为性能瓶颈时,如何利用 异步批量 Sink 将吞吐量从 1w QPS 提升到 10w+?

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
当 Redis 写入成为性能瓶颈时,如何利用 异步批量 Sink 将吞吐量从 1w QPS 提升到 10w+?

引言:从“能跑”到“跑得快”

在上一篇文章中,我们解决了“高可用”问题——通过 Sentinel 让 Flink 任务在 Redis 主从切换时自动恢复。但高可用解决了“活下来”的问题,却未必解决“活得快”的问题。

真实的生产场景往往比我们想象的要残酷。假设你的实时数据流高峰 QPS 达到 5 万、10 万甚至更高,而每个事件都需要写入 Redis。这时你会发现:任务没有崩溃,但反压(Backpressure)越来越严重,吞吐量死活上不去

为什么会这样?我们用官方RedisSink(基于FlinkJedisPoolConfig)逐条写入时,每条数据都要经历一次网络 RTT(Round-Trip Time)。在局域网环境中,单次 RTT 约 0.5~1ms,这意味着单线程极限吞吐只有 1000~2000 QPS。即使开启多个并行度,受限于 Redis 服务端的连接数和处理能力,整体吞吐通常也就在1万~2万 QPS左右。

那么,如何将吞吐量从 1w 提升到 10w+?答案是两个核心技术的组合:异步 I/O + 批量写入(Pipeline)

本文将深入剖析:

  1. 为什么逐条写入 Redis 会成为性能瓶颈——从网络 RTT 到 Redis 服务端处理模型
  2. 异步 I/O 和 Pipeline 各自解决了什么问题
  3. 两种生产级实现方案:基于 Flink Async I/O 和基于 RichSinkFunction 自建批量缓存
  4. 完整的可运行代码、调优参数和常见坑点

一、前置知识:为什么逐条写入这么慢?

1.1 网络 RTT 是最大的隐形杀手

假设你的 Flink 任务和 Redis 部署在同一机房,网络延迟约 0.5ms。逐条写入时,每条数据的处理流程是:

Flink Task -> 获取连接 -> 发送HSET命令 -> 等待Redis响应 -> 归还连接 -> 处理下一条 ↑_______________0.5ms_______________↓

这 0.5ms 的等待时间里,CPU 和网络带宽都在空转。单条数据本身可能只有几百字节,但每次网络往返的“固定开销”远大于数据传输本身的时间。

公式化表达

  • 单条写入耗时 ≈ 网络 RTT(0.5ms)+ Redis 执行时间(~0.05ms)
  • 单线程极限 QPS ≈ 1000ms / 0.55ms ≈1800 QPS

即使开启 10 个并行度,极限也就 1.8 万 QPS——这就是为什么你的任务只能跑到 1w 左右。

1.2 Redis 服务端处理模型的“天花板”

Redis 是单线程处理命令的(指核心事件循环)。这意味着:

  • 无论你有多少个客户端连接,Redis 在同一时刻只能处理一个命令
  • 每个命令的执行时间虽然极短(微秒级),但网络 I/O 和命令解析同样消耗时间

当大量客户端同时涌入时,Redis 的事件循环会出现排队,导致延迟上升。逐条写入放大了这个问题——每个连接每发一条命令就要等一次响应,网络往返次数 = 数据条数。

1.3 两个优化方向

优化方向解决的问题原理
异步 I/O消除 Flink 任务侧的阻塞等待单并行度可同时发起多个未完成的请求
批量写入(Pipeline)减少网络往返次数多条命令合并为一次网络传输

两者结合,理论上可以将吞吐量提升10 倍以上


二、核心剖析:异步 I/O 与 Pipeline 的底层原理

2.1 原理一:Flink Async I/O —— 让“等待”不再阻塞

Flink 在 1.2 版本引入了 Async I/O API。它的核心思想是:

同步模式(MapFunction):

数据1 -> 发请求 -> 阻塞等响应 -> 收到 -> 处理数据2 -> 发请求 -> 阻塞等响应 -> 收到 -> ...

异步模式(AsyncFunction):

数据1 -> 发请求(不等待) 数据2 -> 发请求(不等待) 数据3 -> 发请求(不等待) ...(同时有多个请求在网络上飞行) 响应1回来 -> 处理 响应2回来 -> 处理 响应3回来 -> 处理

单个并行度可以同时发起 N 个未完成的请求(N 由capacity参数控制,默认 100)。这意味着网络等待时间被“重叠”了——在等待响应1的时候,已经在发送请求2、3、4了。

关键参数

  • capacity:最大并发请求数。设置越大吞吐越高,但会增大内存压力和 Redis 服务端压力
  • timeout:请求超时时间。超时未返回的请求会触发异常

注意:Async I/O 更适合读取(维表关联)场景。对于写入(Sink)场景,官方RedisSink并不直接支持 Async I/O。我们需要自建 Sink或使用社区增强版连接器。

2.2 原理二:Redis Pipeline —— 让“多次往返”变成“一次往返”

Redis Pipeline 是 Redis 协议层面的批量处理机制。普通模式下,客户端发送一条命令,必须等到响应后才能发送下一条:

Client: SET key1 value1 Server: +OK Client: SET key2 value2 Server: +OK Client: SET key3 value3 Server: +OK # 3次网络往返

Pipeline 模式下,客户端可以一次性发送多条命令,然后一次性读取所有响应:

Client: SET key1 value1 SET key2 value2 SET key3 value3 ← 三条命令一起发 Server: +OK ← 三条响应一起回 +OK +OK # 1次网络往返

性能提升的数学原理

  • 假设 100 条数据,单条 RTT = 0.5ms
  • 逐条写入:100 × 0.5ms = 50ms
  • Pipeline 批量(100条一批):1 × 0.5ms + 执行时间 ≈ 1ms
  • 提升约 50 倍(理想情况)

实际生产中,受限于网络带宽、Redis 处理能力和批次大小,通常能提升 5~10 倍。有团队在测试中将hincrBy操作从 5w+ ops 提升到了60w+ ops

2.3 两种实现路径对比

实现方式核心机制适用场景复杂度
Flink Async I/O + 逐条写入并发发请求,不阻塞读多写少(维表关联)
RichSinkFunction + Pipeline攒批后批量提交写多读少(Sink 场景)中高
AsyncSink(FLIP-171)Flink 官方异步 Sink API通用 Sink 场景低(需 Flink 1.15+)

对于写入 Redis 的 Sink 场景,最推荐的方案是自定义RichSinkFunction+ Redis Pipeline + 定时 flush


三、手把手实操:两种生产级实现方案

3.1 方案一:基于 Flink Async I/O(适合维表读取场景)

虽然 Async I/O 更适合读取,但如果你需要异步写入 Redis(比如每条数据需要先查 Redis 再决定写什么),这个方案依然适用。

环境依赖(与之前一致):

// build.sbtvalflinkVersion="1.13.6"libraryDependencies++=Seq("org.apache.flink"%%"flink-streaming-scala"%flinkVersion,"redis.clients"%"jedis"%"3.7.0")

核心代码:异步写入 Redis

packageasyncimportorg.apache.flink.streaming.api.scala._importorg.apache.flink.streaming.api.functions.async.{RichAsyncFunction,AsyncFunction}importorg.apache.flink.streaming.api.functions.async.collector.AsyncCollectorimportredis.clients.jedis.{Jedis,JedisPool,JedisPoolConfig}importsource.{Event,ClickSource}importjava.util.concurrent.CompletableFutureimportscala.concurrent.{ExecutionContext,Future}importscala.concurrent.ExecutionContext.Implicits.globalclassAsyncRedisSinkFunction(pool:JedisPool)extendsRichAsyncFunction[Event,Event]{overridedefasyncInvoke(input:Event,collector:AsyncCollector[Event]):Unit={// 使用 CompletableFuture 包装异步操作valfuture=CompletableFuture.supplyAsync(()=>{valjedis=pool.getResourcetry{// 执行 HSET 命令jedis.hset("click",input.user,input.url)input// 返回原数据(或转换为更丰富的类型)}finally{if(jedis!=null)jedis.close()}})// 处理完成回调future.thenAccept(result=>{collector.collect(java.util.Collections.singletonList(result))}).exceptionally(e=>{// 异常处理:可记录日志或发送到死信队列collector.collect(java.util.Collections.emptyList[Event]())null})}overridedeftimeout(input:Event,collector:AsyncCollector[Event]):Unit={// 超时处理collector.collect(java.util.Collections.emptyList[Event]())}}objectAsyncRedisSinkDemo{defmain(args:Array[String]):Unit={valenv=StreamExecutionEnvironment.getExecutionEnvironment env.enableCheckpointing(10000)// 初始化连接池valpoolConfig=newJedisPoolConfig()poolConfig.setMaxTotal(50)poolConfig.setMaxIdle(20)poolConfig.setMinIdle(5)poolConfig.setTestOnBorrow(true)valjedisPool=newJedisPool(poolConfig,"localhost",6379,5000)valdataStream:DataStream[Event]=env.addSource(newClickSource)// 使用 AsyncDataStream 应用异步函数// unorderedWait: 不保证顺序,吞吐更高// orderedWait: 保证顺序,吞吐略低valresultStream=AsyncDataStream.unorderedWait(dataStream,newAsyncRedisSinkFunction(jedisPool),5000,// 超时时间 5 秒java.util.concurrent.TimeUnit.MILLISECONDS,100// capacity: 最大并发请求数)resultStream.print("Written to Redis")env.execute("Async Redis Sink Demo")}}

方案一的局限性

  • AsyncDataStream本质上是流转换算子,不是 Sink。它会产生一个输出流,这在“写入”场景中有些别扭
  • 如果不需要下游继续处理,这种模式会浪费资源
  • 实际吞吐提升约2~3 倍,不如 Pipeline 方案显著

3.2 方案二:RichSinkFunction + Pipeline + 定时 Flush(强烈推荐)

这是生产环境最常用的方案。核心思路:在 Sink 内部维护一个缓冲区,攒够一批数据后使用 Redis Pipeline 一次性提交。

完整代码实现

packagesinkimportorg.apache.flink.streaming.api.scala._importorg.apache.flink.streaming.api.functions.sink.{RichSinkFunction,SinkFunction}importorg.apache.flink.streaming.api.checkpoint.CheckpointedFunctionimportorg.apache.flink.api.common.state.{ListState,ListStateDescriptor}importorg.apache.flink.runtime.state.{FunctionInitializationContext,FunctionSnapshotContext}importorg.apache.flink.streaming.api.functions.sink.SinkFunction.Contextimportorg.apache.flink.util.Preconditionsimportredis.clients.jedis.{Jedis,JedisPool,JedisPoolConfig,Pipeline}importsource.{Event,ClickSource}importscala.collection.mutable.ListBufferimportscala.concurrent.ExecutionContext.Implicits.globalimportscala.concurrent.duration._classBulkRedisSink(host:String,port:Int,batchSize:Int=100,// 批量大小阈值flushIntervalMs:Long=1000// 定时刷新间隔(毫秒))extendsRichSinkFunction[Event]withCheckpointedFunction{// 缓冲区:使用 ListBuffer 存储待写入的数据@transientprivatevarbuffer:ListBuffer[Event]=_@transientprivatevarjedisPool:JedisPool=_@transientprivatevarlastFlushTime:Long=_// Checkpoint 状态:用于故障恢复时重新处理未 flush 的数据@transientprivatevarcheckpointedState:ListState[Event]=_overridedefopen(parameters:org.apache.flink.configuration.Configuration):Unit={// 1. 初始化 Redis 连接池valpoolConfig=newJedisPoolConfig()poolConfig.setMaxTotal(20)poolConfig.setMaxIdle(10)poolConfig.setMinIdle(5)poolConfig.setTestOnBorrow(true)poolConfig.setTestWhileIdle(true)poolConfig.setTimeBetweenEvictionRunsMillis(30000)jedisPool=newJedisPool(poolConfig,host,port,5000)// 2. 初始化缓冲区和最后刷新时间buffer=ListBuffer.empty[Event]lastFlushTime=System.currentTimeMillis()// 3. 注册 ProcessingTimeTimer(定时刷新)valruntimeContext=getRuntimeContext// 注意:此处简化实现,实际可使用 ProcessingTimeService// 或通过单独的调度线程实现定时 flush}overridedefinvoke(value:Event,context:SinkFunction.Context):Unit={// 1. 将数据加入缓冲区buffer.synchronized{buffer+=value}// 2. 检查是否达到批量阈值if(buffer.size>=batchSize){flush()}// 3. 检查是否到达定时刷新时间valnow=System.currentTimeMillis()if(now-lastFlushTime>=flushIntervalMs){flush()}}/** * 核心方法:使用 Pipeline 批量写入 Redis */privatedefflush():Unit={buffer.synchronized{if(buffer.isEmpty)return// 取出当前批次数据(不移除,等写入成功后再清空)valbatch=buffer.toListvarjedis:Jedis=nulltry{jedis=jedisPool.getResourcevalpipeline:Pipeline=jedis.pipelined()// 批量添加命令到 Pipelinebatch.foreach{event=>pipeline.hset("click",event.user,event.url)}// 可选:设置过期时间(如 7 天)// pipeline.expire("click", 604800)// 执行 Pipeline(一次网络往返)pipeline.sync()// 写入成功,清空缓冲区buffer.clear()lastFlushTime=System.currentTimeMillis()}catch{casee:Exception=>// 写入失败:保留缓冲区数据,等待重试// 生产环境可增加重试机制和死信队列println(s"Failed to flush to Redis:${e.getMessage}")throwe}finally{if(jedis!=null)jedis.close()}}}// ===== Checkpoint 相关:保证 Exactly-Once 语义 =====overridedefsnapshotState(context:FunctionSnapshotContext):Unit={// 在 Checkpoint 前强制 flush 所有数据flush()// 将缓冲区数据保存到 Checkpoint 状态中checkpointedState.clear()buffer.synchronized{buffer.foreach(checkpointedState.add)}}overridedefinitializeState(context:FunctionInitializationContext):Unit={valdescriptor=newListStateDescriptor[Event]("bulk-redis-buffer-state",classOf[Event])checkpointedState=context.getOperatorStateStore.getListState(descriptor)// 从 Checkpoint 恢复数据if(context.isRestored){importscala.collection.JavaConverters._ buffer=ListBuffer.empty[Event]checkpointedState.get().asScala.foreach(buffer+=_)println(s"Restored${buffer.size}events from checkpoint")}}overridedefclose():Unit={// 关闭前强制 flush 所有剩余数据flush()if(jedisPool!=null)jedisPool.close()}}objectBulkRedisSinkDemo{defmain(args:Array[String]):Unit={valenv=StreamExecutionEnvironment.getExecutionEnvironment// 必须开启 Checkpoint,否则批量 Sink 在故障时可能丢数据env.enableCheckpointing(10000)valdataStream:DataStream[Event]=env.addSource(newClickSource)// 使用自定义批量 SinkdataStream.addSink(newBulkRedisSink(host="localhost",port=6379,batchSize=100,// 每 100 条 flush 一次flushIntervalMs=1000// 或每 1 秒 flush 一次)).name("Bulk Redis Sink").setParallelism(2)// 建议并行度不要太高env.execute("Bulk Redis Sink Demo")}}

代码要点解析

组件作用关键实现
缓冲区(buffer)攒批ListBuffer[Event],线程安全需手动加锁
Pipeline批量提交jedis.pipelined()+pipeline.sync()
双触发条件兼顾吞吐和延迟buffer.size >= batchSize时间间隔 >= flushIntervalMs
CheckpointedFunction故障恢复snapshotState中 flush + 保存状态
连接池复用连接JedisPool,配置TestOnBorrow确保连接可用

3.3 性能调优参数指南

参数推荐值说明
batchSize100~500批次太小提升不明显,太大会增加延迟和内存压力
flushIntervalMs500~2000低流量时保证数据不积压太久
MaxTotal(连接池)并行度 × 2连接数太少会阻塞,太多会压垮 Redis
并行度1~4Redis 单线程,并行度过高反而增加竞争
setTestOnBorrowtrue防止拿到已断开的连接

一个重要的权衡:批量写入会引入延迟。如果batchSize=100,那么第 1 条数据可能要等第 100 条到达才会被写入,最坏延迟 = 100 条数据的到达间隔。通过flushIntervalMs可以控制最大延迟。


四、进阶思考:当吞吐量继续飙升时

4.1 突破单机 Redis 瓶颈

即使 Pipeline 将单机 Redis 吞吐推到 10w+ QPS,单机 Redis 终究有上限(取决于命令复杂度和网络带宽)。当需要更高吞吐时:

方案原理适用场景
Redis Cluster数据分片,多节点并行吞吐 > 10w QPS
Redis 分区(业务层)按 key 哈希写入不同 Redis 实例不想引入 Cluster 复杂度
异步 Sink(FLIP-171)Flink 1.15+ 官方异步 Sink API希望框架层支持

4.2 Checkpoint 与批量 Sink 的“死锁”风险

这是一个容易被忽略的坑:如果flush()方法在执行 Pipeline 时耗时过长(比如 Redis 响应慢),而 Checkpoint 恰好在此时触发,snapshotState会等待flush()完成,可能导致 Checkpoint 超时。

解决方案

  1. 设置合理的 Pipeline 批次大小,单次 flush 不超过 100ms
  2. snapshotState中设置超时保护
  3. 监控 Checkpoint 耗时,及时调整参数

4.3 数据去重与幂等性

Pipeline 批量写入和逐条写入一样,如果任务故障重启,部分数据可能被重复写入。但由于我们使用的是HSET(幂等命令),重复写入不会造成数据不一致。

如果使用非幂等命令(如LPUSHINCR),则需要:

  1. 在业务层设计去重逻辑(如使用唯一 ID)
  2. 或使用 Flink 的 Exactly-Once Sink(需 Sink 支持两阶段提交)

五、总结

核心要点内容回顾
瓶颈根源逐条写入时,网络 RTT 占用了 90% 以上的时间
异步 I/O让单并行度同时发起多个请求,消除阻塞等待
Redis Pipeline多条命令合并为一次网络传输,减少往返次数
推荐方案RichSinkFunction+ 缓冲区 + Pipeline + 定时 Flush
性能提升从 1w QPS 提升到10w+ QPS,典型场景5~10 倍
关键配置batchSize(100500)、`flushIntervalMs`(5002000ms)、开启 Checkpoint

何时选择哪种方案

场景推荐方案
写入吞吐 < 1w QPS官方RedisSink即可
写入吞吐 1w~5w QPSRichSinkFunction+ Pipeline(本文方案二)
写入吞吐 > 5w QPSPipeline + Redis Cluster + 并行度调优
维表关联(读取)Flink Async I/O(本文方案一)

核心口诀

攒一批,管道送;定时刷,防积压;开 CP,保一致性;调池子,防阻塞。

从“逐条写入”到“批量 Pipeline”,改变的不仅仅是代码,更是对“网络 I/O 是最大瓶颈”这一底层认知的升级。下次当你再遇到 Sink 性能问题时,不妨先问问自己:我的数据是“条条都等”还是“攒批再发”?

下期预告:当 Redis 写入不再是瓶颈后,Flink 任务的反压可能来自哪里?如何系统性地定位和解决 Flink 反压问题?敬请期待。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/10 14:56:03

Spring Boot + Vue 前后端分离实战:高校自习室预约系统开发指南

这次我们来看一个基于 Spring Boot 和 Vue.js 的高校自习室预定系统&#xff0c;项目代号 hx4078。对于高校学生和教务管理者来说&#xff0c;一个稳定、高效、易用的自习室资源管理系统&#xff0c;能直接解决座位难找、资源分配不均、管理混乱的痛点。这个项目就是一个典型的…

作者头像 李华
网站建设 2026/8/10 14:55:35

免费显卡显存测试终极指南:用memtest_vulkan诊断硬件稳定性问题

免费显卡显存测试终极指南&#xff1a;用memtest_vulkan诊断硬件稳定性问题 【免费下载链接】memtest_vulkan Vulkan compute tool for testing video memory stability 项目地址: https://gitcode.com/gh_mirrors/me/memtest_vulkan 你是否遇到过游戏突然闪退、屏幕出现…

作者头像 李华
网站建设 2026/8/10 14:54:41

Tendermint-rs在企业级区块链项目中的应用案例:从理论到实践

Tendermint-rs在企业级区块链项目中的应用案例&#xff1a;从理论到实践 【免费下载链接】tendermint-rs Client libraries for Tendermint/CometBFT in Rust! 项目地址: https://gitcode.com/gh_mirrors/te/tendermint-rs Tendermint-rs是一个用Rust编写的Tendermint/C…

作者头像 李华
网站建设 2026/8/10 14:54:04

如何在Windows通知栏中悄悄完成英语学习的革命性工具

如何在Windows通知栏中悄悄完成英语学习的革命性工具 【免费下载链接】ToastFish 一个利用摸鱼时间背单词的软件。 项目地址: https://gitcode.com/GitHub_Trending/to/ToastFish ToastFish是一款巧妙利用Windows通知栏的智能背单词软件&#xff0c;它将学习过程无缝融入…

作者头像 李华
网站建设 2026/8/10 14:53:40

5分钟掌握猫抓插件:新手必备的浏览器资源嗅探完整指南

5分钟掌握猫抓插件&#xff1a;新手必备的浏览器资源嗅探完整指南 【免费下载链接】cat-catch 猫抓 浏览器资源嗅探扩展 / cat-catch Browser Resource Sniffing Extension 项目地址: https://gitcode.com/GitHub_Trending/ca/cat-catch 你是否经常在网上遇到心仪的视频…

作者头像 李华