Flink连接Elasticache Redis集群报错:无法解析CLUSTER NODES响应cport
问题结论
该问题是低版本Jedis的已知bug+配置类型不匹配共同导致的,不属于业务代码逻辑错误。
- 你使用的端点带
clustercfg前缀,属于Elasticache开启Redis Cluster集群模式的配置端点,本身需要使用Cluster模式客户端连接,但Flink旧版Redis连接器默认依赖的Jedis版本(2.9.0以下)未兼容Redis 4.0+版本CLUSTER NODES响应中ip:port@cport格式的总线端口字段,直接将@后的内容和端口拼接后尝试转整数,触发NumberFormatException。 - 如果你实际创建的是未开启集群模式的普通主从Replication Group,那属于配置类选错:这类实例没有集群槽位逻辑,不需要走Cluster客户端初始化,用Cluster配置连接会主动发送不被支持的
CLUSTER NODES命令,返回内容格式不符合预期触发报错。
解决方案
根据你实际的Replication Group集群模式选择对应方案:
方案1:使用集群模式Replication Group(带clustercfg配置端点)
排除Flink Redis连接器自带的低版本Jedis依赖,引入3.0.0及以上版本的Jedis,该版本已修复cport字段解析问题。
Maven依赖配置参考:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-redis_${scala.binary.version}</artifactId> <version>${flink.version}</version> <exclusions> <exclusion> <groupId>redis.clients</groupId> <artifactId>jedis</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>redis.clients</groupId> <artifactId>jedis</artifactId> <version>3.3.0</version> </dependency>
原有FlinkJedisClusterConfig的代码逻辑不需要修改,升级依赖后即可正常解析集群节点信息。
方案2:使用非集群模式主从Replication Group
替换配置类,不要使用Cluster模式配置,改用FlinkJedisPoolConfig直接连接Replication Group的主节点端点(非clustercfg开头的Primary Endpoint)即可,代码示例:
String primaryEndpoint = "你的Replication Group主节点地址"; int port = 6379; FlinkJedisPoolConfig jedisConfig = new FlinkJedisPoolConfig.Builder() .setHost(primaryEndpoint) .setPort(port) // 按需配置连接超时、密码、连接池大小、数据库编号 .build(); RedisSink redisSink = new RedisSink<>(jedisConfig, new MyRedisMapper());
注意事项
- Kinesis Analytics Flink作业需要和Elasticache实例部署在同一VPC下,关联的安全组需要放通到Elasticache 6379端口的出站规则。
- 如果Elasticache开启了TLS加密、密码认证,需要在Jedis配置中补充开启SSL、设置auth密码的参数。
内容的提问来源于stack exchange,提问作者Jake
相关产品推荐
相关产品推荐

