Docker中无法连接Cassandra,cqlsh报‘Unable to connect to any servers’错误
问题描述
我正在搭建基于PySpark和Docker的容器化批处理数据管道,处理后的数据要存入Cassandra。用docker-compose拉取Spark和Cassandra镜像后,PySpark文件能正常运行,但执行创建Keyspace、表等Cassandra相关操作时出错。进入容器执行cqlsh命令后,出现以下错误:
Connection error: ('Unable to connect to any servers', {
'127.0.0.1:9042': ConnectionRefusedError(111, "Tried connecting to
[('127.0.0.1', 9042)]. Last error: Connection refused")})
使用的Docker命令:
docker compose up -d docker ps docker exec -it container-id cqlsh # 执行此命令后出现错误
我试过多种Cassandra镜像,都出现相同错误,同时还在找如何用Airflow在容器中调度该管道,但没找到有效方案。我的docker-compose配置如下:
version: '3' networks: app-tier: driver: bridge services: spark: image: docker.io/bitnami/spark:3.3 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: - '8080:8080' volumes: - ".:/opt/spark" spark-worker: image: docker.io/bitnami/spark:3.3 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 networks: - app-tier cassandra: image: 'bitnami/cassandra:latest' #image: docker.io/bitnami/cassandra:4.1 #image: cassandra:latest ports: - '7000:7000' - '127.0.0.1:9042:9042' volumes: #- 'cassandra_data:/bitnami' - ".:/opt/cassandra" environment: - CASSANDRA_SEEDS=cassandra - CASSANDRA_PASSWORD_SEEDER=yes - CASSANDRA_PASSWORD=cassandra networks: - app-tier
解决方案
一、修复Cassandra连接问题
1. 调整cqlsh连接命令
Bitnami镜像的Cassandra默认绑定容器内部服务名而非127.0.0.1,进入容器后执行带指定参数的命令:
docker exec -it container-id cqlsh cassandra 9042 -u cassandra -p cassandra
2. 修改docker-compose配置
- 端口映射优化:去掉
127.0.0.1绑定,改为全地址映射,确保容器内外都能访问:ports: - '7000:7000' - '9042:9042' - 数据卷修正:恢复默认数据卷挂载,避免当前目录挂载覆盖容器内部配置和数据:
volumes: - 'cassandra_data:/bitnami' # - ".:/opt/cassandra" # 注释此行,防止干扰容器内部文件 - 添加服务依赖与网络:给Spark服务加入
app-tier网络,并设置依赖Cassandra,确保Cassandra启动后再启动Spark:spark: # ... 原有配置 depends_on: - cassandra networks: - app-tier
3. 验证服务状态
启动容器后查看Cassandra日志,确认服务正常启动:
docker logs cassandra-container-id
日志中出现Startup complete字样说明服务就绪。
二、PySpark连接Cassandra代码调整
PySpark代码中需用Docker网络内的服务名cassandra作为连接地址,而非127.0.0.1:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("CassandraPipeline") \ .config("spark.cassandra.connection.host", "cassandra") \ .config("spark.cassandra.auth.username", "cassandra") \ .config("spark.cassandra.auth.password", "cassandra") \ .getOrCreate() # 创建Keyspace示例 spark.sql("CREATE KEYSPACE IF NOT EXISTS test_keyspace WITH REPLICATION = {'class': 'SimpleStrategy', 'replication_factor': 1}")
三、Airflow容器化调度方案
1. 在docker-compose中添加Airflow服务
在现有配置中追加Airflow相关服务:
services: # ... 现有spark、cassandra服务 postgres: image: postgres:13 environment: - POSTGRES_USER=airflow - POSTGRES_PASSWORD=airflow - POSTGRES_DB=airflow volumes: - postgres_data:/var/lib/postgresql/data networks: - app-tier airflow-webserver: image: apache/airflow:2.6.0 command: webserver environment: - AIRFLOW__CORE__EXECUTOR=LocalExecutor - AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow - AIRFLOW__CORE__LOAD_EXAMPLES=False volumes: - ./dags:/opt/airflow/dags - ./spark_scripts:/opt/airflow/spark_scripts # 挂载PySpark脚本目录 ports: - '8081:8080' depends_on: - postgres - spark - cassandra networks: - app-tier airflow-scheduler: image: apache/airflow:2.6.0 command: scheduler environment: - AIRFLOW__CORE__EXECUTOR=LocalExecutor - AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow - AIRFLOW__CORE__LOAD_EXAMPLES=False volumes: - ./dags:/opt/airflow/dags - ./spark_scripts:/opt/airflow/spark_scripts depends_on: - postgres - airflow-webserver networks: - app-tier volumes: cassandra_data: postgres_data:
2. 编写Airflow DAG调度任务
在./dags目录下创建DAG文件:
from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 1, 1), 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } with DAG( 'spark_cassandra_pipeline', default_args=default_args, description='Run PySpark pipeline to process data and store to Cassandra', schedule_interval=timedelta(days=1), catchup=False, ) as dag: run_spark_job = BashOperator( task_id='run_spark_pipeline', bash_command='spark-submit --master spark://spark:7077 /opt/airflow/spark_scripts/your_pipeline.py', ) run_spark_job
内容的提问来源于stack exchange,提问作者Mikehade
相关产品推荐
相关产品推荐

