Flink 流计算进阶:手写一套通用的 Redis/MySQL /Kafka数据源封装工具类 一、redis篇幅1、解决痛点1、兼容redis集群和单机2、异步IO,通过 AsyncDataStream.unorderedWait 提升吞吐量3、通用性,避免重复造轮子2、代码部分RedisClusterUtil工具类package com.juxin.util.redis; import cn.hutool.core.util.NumberUtil; import cn.hutool.core.util.StrUtil; import com.juxin.util.ConfigUtil; import io.lettuce.core.RedisURI; import io.lettuce.core.cluster.RedisClusterClient; import io.lettuce.core.cluster.api.StatefulRedisClusterConnection; import io.lettuce.core.cluster.api.async.RedisAdvancedClusterAsyncCommands; import io.lettuce.core.cluster.api.sync.RedisAdvancedClusterCommands; import lombok.extern.slf4j.Slf4j; import java.io.Serializable; import java.time.Duration; import java.util.Map; /** * 功能: redis工具类 * 弊端: 同步阻塞,若自创异步会导致还没拿到数据,流已经往下执行了 * */ @Slf4j public class RedisClusterUtil implements Serializable { private static final long serialVersionUID = 1L;//版本校验,Flink忽略类变化,强制读取ck旧数据 private transient RedisClusterClient clusterClient; //跳过不序列化 private transient StatefulRedisClusterConnectionString, String connection; private transient RedisAdvancedClusterCommandsString, String syncCmd; //同步 // 新增:异步命令接口 private transient RedisAdvancedClusterAsyncCommandsString, String asyncCmd; // ✅ 初始化方法,由 Flink 的 open() 调用 public void init() { if (clusterClient != null) {return;} // 已初始化,直接返回 try { String host = ConfigU

本月热点