You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用PySpark连接GKE上Redis集群失败问题排查求助

问题描述

我通过Helm在GKE上部署了Redis集群,并创建了内部负载均衡器,以便Dataproc上运行的PySpark能够访问该集群。操作步骤如下:

  1. 安装Redis,为主节点和从节点添加标签name=myredis:
helm install vani-redis bitnami/redis -n redis --set master.podLabels.name=myredis,replica.podLabels.name=myredis
  1. 内部负载均衡器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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 18:18:11