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

Spark容器无法连接/写入Cassandra容器表的问题求助

问题:Spark集群无法连接Docker网络中的Cassandra服务

我将Kafka Broker、Spark Master/Worker、Cassandra部署在同一Docker confluent 网络下,配置如下:

broker:
    image: confluentinc/cp-server:7.4.0
    hostname: broker
    container_name: broker
    depends_on:
      zookeeper:
        condition: service_healthy
    ports:
      - "9092:9092"
      - "9101:9101"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
      KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_CONFLUENT_BALANCER_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_JMX_PORT: 9101
      KAFKA_JMX_HOSTNAME: localhost
      KAFKA_CONFLUENT_SCHEMA_REGISTRY_URL: http://schema-registry:8081
      CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092
      CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1
      CONFLUENT_METRICS_ENABLE: 'false'
      CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous'
    networks:
      - confluent
    healthcheck:
      test: [ "CMD", "bash", "-c", 'nc -z localhost 9092' ]
      interval: 10s
      timeout: 5s
      retries: 5

  spark-master:
    image: bitnami/spark:latest
    volumes: 
      - ./spark_stream.py:/opt/bitnami/spark/spark_stream.py
    command: bin/spark-class org.apache.spark.deploy.master.Master
    ports:
      - "9090:8080"
      - "7077:7077"
    networks:
      - confluent

  spark-worker:
    image: bitnami/spark:latest
    command: bin/spark-class org.apache.spark.deploy.worker.Worker spark://spark-master:7077
    depends_on:
      - spark-master
    environment:
      SPARK_MODE: worker
      SPARK_WORKER_CORES: 2
      SPARK_WORKER_MEMORY: 1g
      SPARK_MASTER_URL: spark://spark-master:7077
    networks:
      - confluent
    volumes: 
      - ./spark_stream.py:/opt/bitnami/spark/spark_stream.py

  cassandra_db:
    image: cassandra:latest
    container_name: cassandra
    hostname: cassandra
    ports:
      - "9042:9042"
    environment:
      - MAX_HEAP_SIZE=512M
      - HEAP_NEWSIZE=100M
      - CASSANDRA_USERNAME=cassandra
      - CASSANDRA_PASSWORD=cassandra
    networks:
      - confluent


# Define the networks
networks:
  confluent:

代码与错误情况

  1. SparkSession配置:
s_conn = SparkSession.builder \
            .appName('SparkDataStreaming') \
            .config('spark.cassandra.connection.host', 'cassandra') \
            .config("spark.cassandra.connection.port", "9042") \
            .config("spark.cassandra.auth.username", "cassandra") \
            .config("spark.cassandra.auth.password", "cassandra") \
            .getOrCreate()
  1. 本地测试Cassandra连接(成功):
auth_provider = PlainTextAuthProvider(username='cassandra', password='cassandra')
cluster = Cluster(['localhost'], port=9042, auth_provider=auth_provider, connect_timeout=60)
cas_session = cluster.connect()
  1. Spark读取Kafka流:
spark_df = spark_conn.readStream \
            .format('kafka') \
            .option('kafka.bootstrap.servers', 'broker:29092') \
            .option('subscribe', 'users_created') \
            .option('startingOffsets', 'earliest') \
            .load()
  1. 提交命令:
spark-submit --master spark://localhost:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=cassandra --conf spark.cassandra.connection.port=9042 spark_stream.py

错误1:使用cassandra作为主机名时

ERROR CassandraConnectorConf: Unknown host 'cassandra'
java.net.UnknownHostException: cassandra: nodename nor servname provided, or not known

错误2:改用localhost/127.0.0.1时

ERROR TaskSetManager: Task 0 in stage 0.0 failed 4 times; aborting job
24/11/14 01:59:40 ERROR WriteToDataSourceV2Exec: Data source write support MicroBatchWrite[epoch: 0, writer: CassandraBulkWrite(...)] is aborting.

已确认Kafka主题有数据,排除列匹配问题,寻求解决方法。


解决方案

核心原因分析

  1. Spark Driver与Worker的网络隔离:本地执行spark-submit时,Driver运行在宿主机,而Spark Worker在Docker容器的confluent网络内。
    • 用cassandra主机名时,宿主机无法解析Docker网络内的域名,导致Driver找不到Cassandra。
    • 用localhost时,Worker容器内的localhost指向自身,而非宿主机映射的Cassandra服务,Worker无法连接。
  2. Cassandra认证配置缺失:SparkSession中的用户名密码配置可能未传递到Worker节点。

具体修复步骤

1. 统一网络访问方式

方案A:让Driver加入Docker网络

使用同版本Spark镜像在Docker内部提交任务,确保Driver和Worker在同一网络:

docker run --rm --network confluent -v $(pwd):/app bitnami/spark:latest spark-submit --master spark://spark-master:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=cassandra --conf spark.cassandra.connection.port=9042 --conf spark.cassandra.auth.username=cassandra --conf spark.cassandra.auth.password=cassandra /app/spark_stream.py

方案B:用宿主机Docker网络IP配置连接

  • 执行docker network inspect confluent获取宿主机在该网络中的IPv4地址(如172.18.0.1)。
  • 修改SparkSession配置:
    s_conn = SparkSession.builder \
                .appName('SparkDataStreaming') \
                .config('spark.cassandra.connection.host', '172.18.0.1') \
                .config("spark.cassandra.connection.port", "9042") \
                .config("spark.cassandra.auth.username", "cassandra") \
                .config("spark.cassandra.auth.password", "cassandra") \
                .getOrCreate()
    
  • 更新spark-submit命令:
    spark-submit --master spark://localhost:7077 --packages "com.datastax.spark:spark-cassandra-connector_2.12:3.5.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3" --conf spark.cassandra.connection.host=172.18.0.1 --conf spark.cassandra.connection.port=9042 --conf spark.cassandra.auth.username=cassandra --conf spark.cassandra.auth.password=cassandra spark_stream.py
    

2. 确保认证配置全局生效

在spark-submit命令中显式添加认证参数,保证Worker节点能获取到用户名密码配置。

3. 验证网络连通性

在Spark Worker容器内测试Cassandra访问:

docker exec -it <spark-worker-container-id> ping cassandra
docker exec -it <spark-worker-container-id> nc -zv cassandra 9042

若能ping通且端口开放,说明网络无问题,需检查Spark配置或提交方式。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 06:05:57