使用PySpark连接GKE上Redis集群失败问题排查求助
我通过Helm在GKE上部署了Redis集群,并创建了内部负载均衡器,以便Dataproc上运行的PySpark能够访问该集群。操作步骤如下:
- 安装Redis,为主节点和从节点添加标签
name=myredis:
helm install vani-redis bitnami/redis -n redis --set master.podLabels.name=myredis,replica.podLabels.name=myredis
- 内部负载均衡器YAML配置:
apiVersion: v1 kind: Service metadata: name: redis-loadbalancer annotations: networking.gke.io/load-balancer-type: "Internal" spec: type: LoadBalancer externalTrafficPolicy: Cluster selector: name: myredis ports: - name: redis protocol: TCP port: 6379 targetPort: 6379
我可以在Kubernetes的redis命名空间和default命名空间的Pod中通过命令行访问Redis集群,但在Dataproc上运行PySpark程序连接Redis时,抛出如下异常:
Traceback (most recent call last): File "/tmp/a24ebf59-6c7d-4dc9-b759-1eb1e13dfb70/vani-access-redis.py", line 17, in <module> df1.write \ File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 1107, in save File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py", line 1304, in __call__ File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 111, in deco File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o87.save. : redis.clients.jedis.exceptions.JedisConnectionException: Could not get a resource from the pool at redis.clients.jedis.util.Pool.getResource(Pool.java:84) at redis.clients.jedis.JedisPool.getResource(JedisPool.java:377) at com.redislabs.provider.redis.ConnectionPool$.connect(ConnectionPool.scala:35) at com.redislabs.provider.redis.RedisEndpoint.connect(RedisConfig.scala:94) at com.redislabs.provider.redis.RedisConfig.getNonClusterNodes(RedisConfig.scala:276) at com.redislabs.provider.redis.RedisConfig.getNodes(RedisConfig.scala:370) at com.redislabs.provider.redis.RedisConfig.getHosts(RedisConfig.scala:267) at com.redislabs.provider.redis.RedisConfig.<init>(RedisConfig.scala:166) at com.redislabs.provider.redis.RedisConfig$.fromSparkConfAndParameters(RedisConfig.scala:154) at org.apache.spark.sql.redis.RedisSourceRelation.<init>(RedisSourceRelation.scala:34) at org.apache.spark.sql.redis.DefaultSource.createRelation(DefaultSource.scala:21) at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:46) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70) at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68) at org.apache.spark.sql.execution.command.ExecutedCommandExec.doExecute(commands.scala:90) at org.apache.spark.sql.execution.SparkPlan.$anonfun$execute$1(SparkPlan.scala:180) at org.apache.spark.sql.execution.SparkPlan.$anonfun$executeQuery$1(SparkPlan.scala:218) at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:215) at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:176) at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:133) at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:132) at org.apache.spark.sql.DataFrameWriter.$anonfun$runCommand$1(DataFrameWriter.scala:989) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:989) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:438) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:415) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:301) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750) Caused by: redis.clients.jedis.exceptions.JedisConnectionException: Failed to create socket. at redis.clients.jedis.DefaultJedisSocketFactory.createSocket(DefaultJedisSocketFactory.java:110) at redis.clients.jedis.Connection.connect(Connection.java:226) at redis.clients.jedis.BinaryClient.connect(BinaryClient.java:144) at redis.clients.jedis.BinaryJedis.connect(BinaryJedis.java:314) at redis.clients.jedis.BinaryJedis.initializeFromClientConfig(BinaryJedis.java:92) at redis.clients.jedis.BinaryJedis.<init>(BinaryJedis.java:297) at redis.clients.jedis.Jedis.<init>(Jedis.java:169) at redis.clients.jedis.JedisFactory.makeObject(JedisFactory.java:177) at org.apache.commons.pool2.impl.GenericObjectPool.create(GenericObjectPool.java:571) at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:298) at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:223) at redis.clients.jedis.util.Pool.getResource(Pool.java:75) ... 42 more Caused by: java.net.UnknownHostException: vani-redis-master-0.vani-redis-headless.redis.svc.cluster.local at java.net.AbstractPlainSocketImpl.connect(AbstractPlainSocketImpl.java:184) at java.net.SocksSocketImpl.connect(SocksSocketImpl.java:392) at java.net.Socket.connect(Socket.java:607) at redis.clients.jedis.DefaultJedisSocketFactory.createSocket(DefaultJedisSocketFactory.java:80) ... 53 more
redis命名空间下的资源列表如下:
Karans-MacBook-Pro:redis-on-gke karanalang$ kc get all -n redis NAME READY STATUS RESTARTS AGE pod/vani-redis-master-0 1/1 Running 0 11m pod/vani-redis-replicas-0 1/1 Running 0 11m pod/vani-redis-replicas-1 1/1 Running 0 10m pod/vani-redis-replicas-2 1/1 Running 0 10m NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE service/redis-loadbalancer LoadBalancer 10.228.87.155 xx.xxx.x.xxx 6379:32689/TCP 6m5s service/vani-redis-headless ClusterIP None <none> 6379/TCP 11m service/vani-redis-master ClusterIP 10.228.88.78 <none> 6379/TCP 11m service/vani-redis-replicas ClusterIP 10.228.243.253 <none> 6379/TCP 11m NAME READY AGE statefulset.apps/vani-redis-master 1/1 11m statefulset.apps/vani-redis-replicas 3/3 11m
我的PySpark程序中传入了负载均衡器IP和密码:
log_validation_table = 'redis_table' data = [(1,'k1','address1'),(2,'k2','address2')] df1 = spark.createDataFrame(data, ('id','name', 'address')) redis_table = "redis_table" df1.write \ .format("org.apache.spark.sql.redis") \ .option("table", log_validation_table) \ .option("key.column", "id") \ .option("host", redis_host) \ .option("auth", passwd) \ .option("port", 6379) \ .option("ttl", 2592000) \ .mode("append") \ .save()
请问问题可能出在哪里?如何调试和解决?
核心原因
异常中的UnknownHostException指向Redis集群的内部域名vani-redis-master-0.vani-redis-headless.redis.svc.cluster.local——这是Kubernetes集群内部的DNS解析地址,Dataproc集群不在同一个K8s网络环境中,无法解析该域名。
本质问题是:Spark Redis连接器在连接到Redis主节点后,会从Redis的INFO命令结果中获取集群节点的内部域名,尝试直接连接这些节点,但Dataproc无法访问这些内部域名,最终导致连接失败。
调试步骤
- 验证网络连通性:在Dataproc集群节点上,用
ping <LB-IP>测试负载均衡器IP是否可达,用telnet <LB-IP> 6379测试端口是否开放。 - 检查Redis节点返回地址:在K8s内部Pod中执行
redis-cli -h <LB-IP> -a <密码> INFO replication,查看master_host等字段的值,确认返回的是内部域名还是负载均衡器IP。 - 测试手动连接:在Dataproc节点上用
redis-cli -h <LB-IP> -a <密码>手动连接Redis,执行SET test 123等简单命令,验证基础连接是否正常。
解决方法
方法1:配置Redis集群对外暴露负载均衡器IP
Bitnami Redis支持通过Helm配置让节点对外报告负载均衡器IP,而非内部域名。执行以下命令更新配置:
helm upgrade vani-redis bitnami/redis -n redis \ --set master.podLabels.name=myredis,replica.podLabels.name=myredis \ --set master.replicaAnnounceIp=<你的内部LB-IP> \ --set replica.replicaAnnounceIp=<你的内部LB-IP>
更新完成后重启Redis集群,再检查Redis返回的节点地址是否已替换为负载均衡器IP。
方法2:修改Spark Redis连接器配置,强制使用指定主机
在PySpark代码中添加参数,强制连接器仅使用负载均衡器IP,不自动发现集群节点:
df1.write \ .format("org.apache.spark.sql.redis") \ .option("table", log_validation_table) \ .option("key.column", "id") \ .option("host", redis_host) \ .option("auth", passwd) \ .option("port", 6379) \ .option("ttl", 2592000) \ .option("redis.cluster.mode", "standalone") \ .option("redis.cluster.nodes", f"{redis_host}:6379") \ .mode("append") \ .save()
注意:该方法适用于不需要Redis集群读写分离的场景,所有请求都会通过负载均衡器转发。
方法3:打通Dataproc与GKE的网络
如果Dataproc和GKE不在同一个VPC,可配置VPC peering,让Dataproc能够解析K8s内部DNS域名。此配置复杂度较高,适合长期跨集群网络打通需求。
内容的提问来源于stack exchange,提问作者Karan Alang

