Spark-Cassandra连接器对接K8s Cassandra集群的配置及性能优化求助
问题:Spark-Cassandra连接器写入K8s部署的Cassandra性能极差
问题详情
我们部署在Kubernetes上的Cassandra集群,使用Spark-Cassandra连接器写入时性能极低:
- 写入数据:13亿唯一键的DataFrame(约30GB)
- Spark配置:16个执行器,每个4核、16GB内存
- Cassandra集群:5节点,复制因子=2
- 表结构:
CREATE TABLE <tablename> (hashed_id text PRIMARY KEY, timestamp1 bigint, timestamp2 bigint)
- 写入耗时:约8小时
写入代码示例:
df .write .format("org.apache.spark.sql.cassandra") .mode("overwrite") .option("confirm.truncate", "true") .options(table=tablename, keyspace=cassandra_keyspace) .save()
环境配置
- Cassandra 4.0通过K8ssandra Operator 1.6部署在K8s,前置Traefik Ingress(无TLS)
- Spark 3.2部署在裸金属服务器,ETL基于PySpark,连接器版本:
spark-cassandra-connector_2.12-3.2.0
推测原因
Spark-Cassandra连接器通过Ingress地址连接Cassandra后,获取到的集群节点地址为K8s内部IP,裸金属Spark集群无法直接访问这些内部IP,导致所有写入流量只能走Ingress单点,无法利用Cassandra集群的全部节点能力,形成性能瓶颈。
解决方案及配置示例
1. 调整K8ssandra配置,让Cassandra节点暴露外部可访问地址
修改CassandraDatacenter资源配置,让每个Cassandra节点广播裸金属服务器可访问的IP+端口,同时通过NodePort服务暴露每个节点的CQL端口(9042):
apiVersion: cassandra.datastax.com/v1beta1 kind: CassandraDatacenter metadata: name: dc1 spec: clusterName: cassandra-cluster serverType: cassandra serverVersion: "4.0.11" managementApiAuth: insecure: {} size: 5 storageConfig: cassandraDataVolumeClaimSpec: storageClassName: standard accessModes: - ReadWriteOnce resources: requests: storage: 100Gi podTemplateSpec: spec: containers: - name: cassandra env: # 广播地址设置为K8s节点的外部IP(裸金属可访问) - name: CASSANDRA_BROADCAST_ADDRESS valueFrom: fieldRef: fieldPath: status.hostIP # 监听地址为Pod内部IP - name: CASSANDRA_LISTEN_ADDRESS valueFrom: fieldRef: fieldPath: status.podIP # 为每个Cassandra节点创建NodePort服务(或使用LoadBalancer,根据环境选择) services: - name: cassandra-external type: NodePort ports: - port: 9042 targetPort: 9042 # 可指定固定NodePort,或让K8s自动分配,记录每个节点的端口
注意:如果使用NodePort,需要为每个Cassandra节点单独分配服务(或使用Headless服务结合NodePort),确保每个节点的9042端口映射到裸金属的不同端口,比如节点1映射30001,节点2映射30002等。
2. 配置Spark-Cassandra连接器,指定所有Cassandra节点的外部地址
在PySpark会话中配置连接器直接连接所有Cassandra节点的外部地址,同时调整写入优化参数:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("CassandraHighPerfWrite") \ # 指定所有Cassandra节点的外部IP+NodePort .config("spark.cassandra.connection.host", "node1-ip:30001,node2-ip:30002,node3-ip:30003,node4-ip:30004,node5-ip:30005") \ .config("spark.cassandra.connection.port", "9042") \ # 调整写入批次大小(根据数据大小调整) .config("spark.cassandra.output.batch.size.rows", "2000") \ # 每个执行器的并发写入数 .config("spark.cassandra.output.concurrent.writes", "15") \ # 保持连接活跃,避免频繁重连 .config("spark.cassandra.connection.keep_alive_ms", "60000") \ # 调整Shuffle分区数,匹配Cassandra集群的并行能力(建议节点数*核数*2) .config("spark.sql.shuffle.partitions", "40") \ .getOrCreate() # 执行写入 df.write \ .format("org.apache.spark.sql.cassandra") \ .mode("overwrite") \ .option("confirm.truncate", "true") \ .options(table=tablename, keyspace=cassandra_keyspace) \ .save()
3. 额外优化建议
- 确保Spark执行器的内存分配合理,预留足够内存给Cassandra连接器的写入缓存
- 对DataFrame进行预分区,让分区键与Cassandra的主键(hashed_id)对齐,减少数据 shuffle
- 根据业务需求调整
spark.cassandra.output.consistency.level,比如设置为LOCAL_QUORUM平衡一致性与写入性能
内容的提问来源于stack exchange,提问作者Yiftizur
相关产品推荐
相关产品推荐

