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

Docker容器中Spark与Cassandra连接失败问题求助

问题排查:Spark-Cassandra连接器数据源找不到错误

问题描述

通过docker-compose部署Spark(master/worker)与Cassandra容器,Cassandra已导入数据。将spark-cassandra-connector_2.12-latest_version.jar放入Spark容器jars目录后,运行PySpark读取脚本时出现[DATA_SOURCE_NOT_FOUND]错误,提示找不到org.apache.spark.sql.cassandra数据源。Spark Hello World程序可正常运行,所有容器处于同一Docker网络。

环境配置

docker-compose配置

spark:
    image: docker.io/bitnami/spark:3.4
    container_name: spark_master
    environment:
      - SPARK_MODE=master
      - SPARK_RPC_AUTHENTICATION_ENABLED=no
      - SPARK_RPC_ENCRYPTION_ENABLED=no
      - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
      - SPARK_SSL_ENABLED=no
      - SPARK_USER=spark
    ports:
      - '8081:8081'
  spark-worker:
    image: docker.io/bitnami/spark:3.4
    container_name: spark_worker
    environment:
      - SPARK_MODE=worker
      - SPARK_MASTER_URL=spark://spark:7077
      - SPARK_WORKER_MEMORY=1G
      - SPARK_WORKER_CORES=1
      - SPARK_RPC_AUTHENTICATION_ENABLED=no
      - SPARK_RPC_ENCRYPTION_ENABLED=no
      - SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
      - SPARK_SSL_ENABLED=no
      - SPARK_USER=spark
cassandra:
    container_name: cassandra-container
    image: cassandra:latest
    ports:
      - "7000:7000"
      - "7001:7001"
      - "7199:7199"
      - "9042:9042"
      - "9160:9160"
    volumes:
        - cassandra-data:/var/lib/cassandra
    restart: always

PySpark连接脚本

from pyspark.sql import SparkSession

# Define your Cassandra connection settings
cassandra_host = "127.0.0.1"  # Replace with the actual IP address of your Cassandra container
cassandra_port = "9042"  # Default Cassandra port
cassandra_keyspace = "mykey_space"  # Replace with your Cassandra keyspace
cassandra_table = "test_table"  # Replace with your Cassandra table

# Create a Spark session
spark = SparkSession.builder \
    .appName("CassandraConnectionTest") \
    .config("spark.cassandra.connection.host", cassandra_host) \
    .config("spark.cassandra.connection.port", cassandra_port) \
    .getOrCreate()

# Test the Cassandra connection
try:
    df = spark.read \
        .format("org.apache.spark.sql.cassandra") \
        .options(table=cassandra_table,keyspace= cassandra_keyspace) \
        .load()
    df.show(5)
    print("Cassandra connection test successful.")
except Exception as e:
    print(f"Error: {str(e)}")
    print("Cassandra connection test failed.")

# Stop the Spark session
spark.stop()

运行命令

docker exec -it spark_worker /bin/bash -c "spark-submit --jars /opt/bitnami/spark/jars/spark-cassandra-connector_2.12-latest_version.jar /opt/bitnami/spark/spark-script/cassandra_spark.py"

Docker网络信息

[{
    "Name": "project_default",
    "Id": "c28118c6ba0511ec9f54a97a6d7b17113a75995d4bfed6d18fc80cf3616acb03",
    "Created": "2023-09-18T12:19:08.729378699Z",
    "Scope": "local",
    "Driver": "bridge",
    "EnableIPv6": false,
    "IPAM": {
        "Driver": "default",
        "Options": null,
        "Config": [{
            "Subnet": "172.20.0.0/16",
            "Gateway": "172.20.0.1"
        }]
    },
    "Internal": false,
    "Attachable": false,
    "Ingress": false,
    "ConfigFrom": {
        "Network": ""
    },
    "ConfigOnly": false,
    "Containers": {
        "89396ef426b166e4d293011ef2c8c8ff60007d2baeba2734ef2e91a0cf7f60ec": {
            "Name": "cassandra-container",
            "EndpointID": "5a3804e47d659bcdf19a44477e85fc7c250b52b0375d342984c7d0033ff1b6cc",
            "MacAddress": "02:42:ac:14:00:04",
            "IPv4Address": "172.20.0.4/16",
            "IPv6Address": ""
        },
        "95cdbc38d6717b0318e81e33dfe990dc5da461d22f60976d132559d568c9365f": {
            "Name": "spark_worker",
            "EndpointID": "cc6b69c4c2f2de02965efc78a53dc92d2f06e4e2a0d1e9dd29fb161308032a26",
            "MacAddress": "02:42:ac:14:00:03",
            "IPv4Address": "172.20.0.3/16",
            "IPv6Address": ""
        },
        "e4864bdc0ff72f0185aceaf6f4844ef75ec45f6673cc776529eb8ca2f67943c6": {
            "Name": "spark_master",
            "EndpointID": "83c4b290a2730cec39849189b326513b33b70042832dacf9c5b47b2d92f90351",
            "MacAddress": "02:42:ac:14:00:02",
            "IPv4Address": "172.20.0.2/16",
            "IPv6Address": ""
        }
    },
    "Options": {},
    "Labels": {
        "com.docker.compose.network": "default",
        "com.docker.compose.project": "project",
        "com.docker.compose.version": "2.18.1"
    }
}]

错误日志

23/09/17 17:45:10 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
23/09/17 17:45:11 INFO SparkContext: Running Spark version 3.4.1
...
Error: An error occurred while calling o34.load.
: org.apache.spark.SparkClassNotFoundException: [DATA_SOURCE_NOT_FOUND] Failed to find the data source: org.apache.spark.sql.cassandra. Please find packages at `https://spark.apache.org/third-party-projects.html`.
...
Caused by: java.lang.ClassNotFoundException: org.apache.spark.sql.cassandra.DefaultSource
...

排查与解决方案

1. 确认连接器版本兼容性

Spark 3.4.1对应Scala 2.12,必须使用匹配版本的spark-cassandra-connector,正确版本为3.4.0,文件名应为spark-cassandra-connector_2.12-3.4.0.jar。避免使用latest_version这类模糊命名,防止下载到不兼容版本。

2. 检查jar包路径与权限

  • 执行以下命令确认jar包存在于指定目录:
    docker exec spark_worker ls /opt/bitnami/spark/jars/ | grep cassandra
    
  • 若jar包存在但权限异常,修正权限:
    docker exec spark_worker chown spark:spark /opt/bitnami/spark/jars/spark-cassandra-connector_2.12-3.4.0.jar
    
  • 或者直接省略--jars参数,将jar包放在Spark默认jars目录,Spark会自动加载。

3. 修正Cassandra连接地址

容器内部无法通过127.0.0.1访问Cassandra,需替换为Docker网络内的容器名cassandra-container或内网IP172.20.0.4,修改PySpark脚本中的cassandra_host:

cassandra_host = "cassandra-container"

4. 完善SparkSession配置

在构建SparkSession时添加连接器包配置,让Spark自动管理依赖(需容器联网):

spark = SparkSession.builder \
    .appName("CassandraConnectionTest") \
    .config("spark.cassandra.connection.host", "cassandra-container") \
    .config("spark.jars.packages", "com.datastax.spark:spark-cassandra-connector_2.12:3.4.0") \
    .getOrCreate()

此时运行命令无需指定--jars参数:

docker exec -it spark_worker /bin/bash -c "spark-submit /opt/bitnami/spark/spark-script/cassandra_spark.py"

内容的提问来源于stack exchange,提问作者Amine SAIHI

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 13:32:02