引言:从“能跑”到“跑得快”
在上一篇文章中,我们解决了“高可用”问题——通过 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)。
本文将深入剖析:
- 为什么逐条写入 Redis 会成为性能瓶颈——从网络 RTT 到 Redis 服务端处理模型
- 异步 I/O 和 Pipeline 各自解决了什么问题
- 两种生产级实现方案:基于 Flink Async I/O 和基于 RichSinkFunction 自建批量缓存
- 完整的可运行代码、调优参数和常见坑点
一、前置知识:为什么逐条写入这么慢?
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 性能调优参数指南
| 参数 | 推荐值 | 说明 |
|---|---|---|
batchSize | 100~500 | 批次太小提升不明显,太大会增加延迟和内存压力 |
flushIntervalMs | 500~2000 | 低流量时保证数据不积压太久 |
MaxTotal(连接池) | 并行度 × 2 | 连接数太少会阻塞,太多会压垮 Redis |
并行度 | 1~4 | Redis 单线程,并行度过高反而增加竞争 |
setTestOnBorrow | true | 防止拿到已断开的连接 |
一个重要的权衡:批量写入会引入延迟。如果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 超时。
解决方案:
- 设置合理的 Pipeline 批次大小,单次 flush 不超过 100ms
- 在
snapshotState中设置超时保护 - 监控 Checkpoint 耗时,及时调整参数
4.3 数据去重与幂等性
Pipeline 批量写入和逐条写入一样,如果任务故障重启,部分数据可能被重复写入。但由于我们使用的是HSET(幂等命令),重复写入不会造成数据不一致。
如果使用非幂等命令(如LPUSH、INCR),则需要:
- 在业务层设计去重逻辑(如使用唯一 ID)
- 或使用 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 QPS | RichSinkFunction+ Pipeline(本文方案二) |
| 写入吞吐 > 5w QPS | Pipeline + Redis Cluster + 并行度调优 |
| 维表关联(读取) | Flink Async I/O(本文方案一) |
核心口诀:
攒一批,管道送;定时刷,防积压;开 CP,保一致性;调池子,防阻塞。
从“逐条写入”到“批量 Pipeline”,改变的不仅仅是代码,更是对“网络 I/O 是最大瓶颈”这一底层认知的升级。下次当你再遇到 Sink 性能问题时,不妨先问问自己:我的数据是“条条都等”还是“攒批再发”?
下期预告:当 Redis 写入不再是瓶颈后,Flink 任务的反压可能来自哪里?如何系统性地定位和解决 Flink 反压问题?敬请期待。