尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

当 Redis 集群发生主从切换(Failover)时,Flink 任务会崩溃吗?如何利用 Sentinel 实现高可用?

当 Redis 集群发生主从切换(Failover)时,Flink 任务会崩溃吗?如何利用 Sentinel 实现高可用? 引言从“能用”到“高可用”在上一篇文章中我们基于FlinkJedisPoolConfig构建了一个可运行的 Flink Redis Sink。但生产环境从来不是“能跑就行”——当 Redis 单节点宕机时整个 Flink 任务会直接崩溃因为连接池配置的localhost:6379已经不可用了。那么问题来了如果部署了 Redis 主从Sentinel 高可用架构Flink 任务能在主从切换时自动恢复吗答案是能但有前提——你必须使用FlinkJedisSentinelConfig替代FlinkJedisPoolConfig并且正确配置重试机制。否则任务依然会崩溃。本文将深入剖析Redis Sentinel 的故障转移机制及其对 Flink 任务的影响Flink Redis Connector 在 Sentinel 模式下的连接管理原理一份可直接运行的 Sentinel 高可用配置代码主从切换期间的“数据黑洞”风险与应对策略生产环境必做的 5 项高可用加固措施一、前置知识Redis Sentinel 如何实现高可用1.1 主从复制 哨兵 自动故障转移Redis Sentinel哨兵是 Redis 官方提供的高可用解决方案。它的核心工作流程如下阶段哨兵行为耗时监控每 1 秒向主从节点发送PING检测存活状态持续主观下线SDOWN单个哨兵发现主节点无响应标记为 SDOWN约 3 秒客观下线ODOWN多数哨兵quorum确认主节点不可用触发故障转移取决于 quorum 配置Leader 选举哨兵集群通过 Raft 协议选举一个 Leader 执行故障转移数秒主从切换Leader 将一个从节点提升为新主节点其他从节点切换复制目标数秒客户端感知哨兵将新主节点信息通知客户端通过SENTINEL GET-MASTER-ADDR-BY-NAME即时关键数据在 3 节点 Sentinel 集群中需至少 2 个节点确认故障才会触发 ODOWN避免网络分区导致的误切换。整个故障转移过程通常在10~30 秒内完成。1.2 Flink 任务在故障转移期间会发生什么假设你的 Flink 任务正在向 Redis 主节点master-1:6379写入数据此时主节点宕机阶段一故障发生 → 哨兵检测0~5 秒Flink 任务尝试写入 Redis但连接已断开。Jedis 客户端抛出异常如JedisConnectionException或SocketTimeoutException。此时 Flink 任务会报错但尚未崩溃——错误会被 Flink 的容错机制捕获。阶段二哨兵选举 主从切换5~15 秒哨兵集群正在选举新主节点Redis 服务暂时不可用。Flink 任务持续重试写入不断报错。如果重试策略配置不当任务会在多次失败后崩溃。阶段三新主节点上线 客户端感知15~30 秒新主节点原slave-1已提升为master-2。Sentinel 客户端JedisSentinelPool通过哨兵获取新主节点地址。如果使用了FlinkJedisSentinelConfig连接池会自动发现新主节点并重建连接。任务恢复写入数据继续流向新主节点。结论FlinkJedisPoolConfig直连单节点→任务崩溃FlinkJedisSentinelConfig通过哨兵获取主节点→任务短暂报错后自动恢复。二、核心剖析FlinkJedisSentinelConfig 的工作原理2.1 三种配置类的本质区别Flink Redis Connector 提供了三种配置类对应三种 Redis 部署模式配置类适用场景主从切换时行为FlinkJedisPoolConfig单机 Redis❌ 连接失效任务崩溃FlinkJedisClusterConfigRedis Cluster 集群模式⚠️ 客户端可感知拓扑变化但对主从切换支持有限FlinkJedisSentinelConfigRedis Sentinel 高可用模式✅ 自动发现新主节点连接自动恢复2.2 Sentinel 模式下的连接获取流程当RedisSink使用FlinkJedisSentinelConfig时内部连接管理流程如下1. RedisSink.open() 被调用 ↓ 2. RedisCommandsContainerBuilder.build(jedisSentinelConfig) ↓ 3. 创建 JedisSentinelPool内部持有哨兵地址列表 ↓ 4. JedisSentinelPool 调用 SENTINEL GET-MASTER-ADDR-BY-NAME master ↓ 5. 获取当前主节点 IP 端口 ↓ 6. 建立到主节点的连接池 ↓ 7. 每条数据写入时从连接池借用 Jedis 连接 → 执行 HSET → 归还关键差异JedisSentinelPool不会在初始化时“固定”一个主节点地址。每次获取连接时它都会先向哨兵查询当前主节点再建立连接。这意味着即使主节点发生切换下一次获取连接时就能拿到新主节点地址。2.3 一个容易被忽略的坑连接池缓存JedisSentinelPool内部有一个master字段缓存了主节点地址。如果故障转移发生在两次连接获取之间缓存的地址可能已经失效。Jedis 的处理方式是当使用缓存地址连接失败时会重新向哨兵查询并更新缓存。但这个过程需要时间且可能抛出异常。因此Flink 任务在切换瞬间仍可能出现短暂报错——这是正常现象只要配置了重试机制任务就能自愈。三、手把手实操从PoolConfig升级到SentinelConfig3.1 环境准备搭建 Redis Sentinel 集群最低配置1 个主节点 1 个从节点 3 个哨兵生产环境哨兵至少 3 个# 以 Docker Compose 为例快速验证用version:3services: redis-master: image: redis:7 command: redis-server--port6379--appendonlyyesports: -6379:6379redis-slave: image: redis:7 command: redis-server--port6380--slaveofredis-master6379--appendonlyyesports: -6380:6380sentinel-1: image: redis:7 command: redis-sentinel /usr/local/etc/redis/sentinel.conf volumes: - ./sentinel.conf:/usr/local/etc/redis/sentinel.conf ports: -26379:26379# sentinel-2, sentinel-3 同理...sentinel.conf核心配置port 26379 sentinel monitor mymaster redis-master 6379 2 sentinel down-after-milliseconds mymaster 5000 sentinel failover-timeout mymaster 600003.2 升级后的 Scala 代码完整可运行packagesinkimportorg.apache.flink.streaming.api.scala._importorg.apache.flink.streaming.connectors.redis.RedisSinkimportorg.apache.flink.streaming.connectors.redis.common.config.FlinkJedisSentinelConfigimportorg.apache.flink.streaming.connectors.redis.common.mapper.{RedisCommand,RedisCommandDescription,RedisMapper}importsource.{ClickSource,Event}importjava.util.{HashSetJHashSet}importscala.collection.JavaConverters._objectsinkToRedisSentinel{defmain(args:Array[String]):Unit{valenvStreamExecutionEnvironment.getExecutionEnvironment env.enableCheckpointing(10000)valdataStream:DataStream[Event]env.addSource(newClickSource)dataStream.print(Input from Source)// 1. 配置哨兵地址至少 3 个valsentinelsnewJHashSet[String]()sentinels.add(sentinel-1:26379)sentinels.add(sentinel-2:26379)sentinels.add(sentinel-3:26379)// 2. 构建 Sentinel 配置核心valconf:FlinkJedisSentinelConfignewFlinkJedisSentinelConfig.Builder().setMasterName(mymaster)// 必须与 sentinel.conf 中的名称一致.setSentinels(sentinels)// 哨兵地址列表.setConnectionTimeout(5000)// 连接超时 5 秒.setSoTimeout(5000)// Socket 超时 5 秒.setMaxTotal(20)// 最大连接数.setMaxIdle(10)// 最大空闲连接.setMinIdle(5)// 最小空闲连接.setTestOnBorrow(true)// 借用时检查连接可用性.setTestWhileIdle(true)// 空闲时检查连接可用性.build()// 3. 添加 Redis SinkMapper 逻辑与之前完全一致valredisSinknewRedisSink[Event](conf,newRedisMapper[Event]{overridedefgetCommandDescription:RedisCommandDescriptionnewRedisCommandDescription(RedisCommand.HSET,click)overridedefgetKeyFromData(t:Event):Stringt.useroverridedefgetValueFromData(t:Event):Stringt.url})dataStream.addSink(redisSink).name(Redis Sentinel Sink).setParallelism(1)env.execute(Flink Redis Sentinel Job)}}代码变化对比原版PoolConfig新版SentinelConfig.setHost(localhost).setSentinels(sentinels).setMasterName(mymaster)直连单节点通过哨兵动态发现主节点主从切换 → 任务崩溃主从切换 → 自动重连3.3 验证 Sentinel 高可用效果Step 1启动 Flink 任务确认数据正常写入 Redis。redis-cli-hlocalhost-p6379HGETALL click# 正常返回数据Step 2模拟主节点宕机。dockerstop redis-masterStep 3观察 Flink 任务日志。# 预期会看到类似以下日志非精确取决于 Jedis 版本 WARN JedisSentinelPool - master mymaster is down, trying to discover new master WARN JedisSentinelPool - new master found: redis-slave:6380 INFO RedisSink - reconnected to Redis successfullyStep 4验证数据仍在写入新主节点。redis-cli-hlocalhost-p6380HGETALL click# 数据持续增长说明故障转移成功四、进阶思考主从切换期间的“数据黑洞”4.1 异步复制导致的数据丢失风险Redis 主从复制是异步的。这意味着主节点收到写入请求 → 返回OK给客户端 →然后才异步同步到从节点。如果在同步完成之前主节点宕机这部分数据永远丢失了。这就是所谓的“数据黑洞”——在故障转移期间部分已确认写入的数据可能永久消失。对 Flink 任务的影响Flink 的 Checkpoint 机制认为数据已成功写入因为RedisSink返回了成功。但实际上数据并未复制到从节点主节点宕机后数据丢失。Flink 的 Exactly-Once 语义无法覆盖这种场景因为数据丢失发生在 Redis 内部Flink 无法感知。4.2 应对策略策略实现方式优缺点开启 AOF 持久化appendonly yesappendfsync always数据更安全但性能下降明显使用 WAIT 命令主节点写入后等待从节点确认牺牲可用性换取一致性业务层容忍丢失接受故障转移期间少量数据丢失大多数实时场景可接受双写 去重同时写入两个 Redis 集群消费端去重成本翻倍复杂度高生产建议对于实时点击流这类非关键数据可以容忍少量丢失对于交易数据则不应依赖 Redis 作为唯一存储而应使用支持事务的数据库。五、生产环境必做的 5 项高可用加固5.1 配置 Flink 重启策略// 在 env 创建后立即配置env.setRestartStrategy(RestartStrategies.fixedDelayRestart(10,// 最多重试 10 次Time.seconds(30)// 每次重试间隔 30 秒))为什么重要故障转移期间 Redis 可能不可用 10~30 秒如果没有重启策略任务会在第一次报错时就失败。5.2 设置合理的超时时间.setConnectionTimeout(10000)// 连接超时 10 秒故障转移期间可能较长.setSoTimeout(10000)// 读取超时 10 秒为什么重要故障转移期间连接建立可能比平时慢过短的超时会导致误判。5.3 开启连接池健康检查.setTestOnBorrow(true).setTestWhileIdle(true).setTimeBetweenEvictionRunsMillis(30000)// 每 30 秒检查一次空闲连接为什么重要主从切换后旧主节点的连接已失效。开启TestOnBorrow可以在借用连接时执行PING检测确保不会拿到死连接。5.4 监控告警监控项告警阈值意义Flink 任务重启次数 3 次/小时可能存在持续性故障Redis 连接异常日志任何JedisConnectionException主从切换或网络问题Redis 主从延迟 1 秒异步复制积压可能丢数据哨兵集群健康任意哨兵不可达哨兵集群本身出问题5.5 避免雪崩连接池预热与限流故障恢复后所有并行子任务可能同时尝试重新连接 Redis瞬间产生大量连接请求。应对// 方式一限制 Sink 并行度.setParallelism(1)// 方式二在连接池配置中限制最大连接数.setMaxTotal(10)// 不要设置过大避免恢复时打爆 Redis六、总结问题答案Redis 主从切换时 Flink 任务会崩溃吗使用FlinkJedisPoolConfig→会崩溃使用FlinkJedisSentinelConfig→不会崩溃会自动恢复。故障转移期间会发生什么Flink 任务会短暂报错连接超时但配合重启策略和 Sentinel 自动重连任务可在 10~30 秒内自愈。数据会丢失吗可能。Redis 异步复制导致主从切换时存在“数据黑洞”需根据业务重要性决定是否容忍。生产环境还需要做什么配置重启策略、合理超时、连接池健康检查、监控告警、防止恢复时的连接风暴。核心要点回顾配置类升级从FlinkJedisPoolConfig切换到FlinkJedisSentinelConfig是让 Flink 任务在 Redis 主从切换时“活下来”的第一步。重启策略是保底即使 Sentinel 自动重连故障转移期间的短暂不可用仍可能触发 Flink 报错。fixedDelayRestart是必选项。数据一致性是上限Redis 的 AP 特性决定了它在故障转移时无法保证数据零丢失。关键数据请使用支持事务的存储系统。监控是最后的防线没有监控的高可用是“伪高可用”——你永远不会知道故障发生过直到数据出问题。下期预告当 Redis 写入成为性能瓶颈时如何利用异步批量 Sink将吞吐量从 1w QPS 提升到 10w敬请期待。
返回列表