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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 18:47:08