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

在Apache Airflow中使用Spark Submit Operator对接远程Cassandra服务器

问题排查与解决方案:Docker中Airflow通过Spark写入远程Cassandra失败

针对能读取Cassandra数据但无法写入、调用df.write.save()报错An error occured while calling o41.save.的问题,可从以下方向排查解决:

1. 修正Spark Cassandra Connector的配置参数

你的Spark配置存在两个关键问题:

  • spark.jars.packages用于指定Maven仓库依赖坐标,本地Jar包路径应使用spark.jars参数。
  • 配置中Jar文件名后缀错误(.jars应为.jar)。

修正后的SparkSession配置示例:

spark = SparkSession \
        .builder \
        .master("local[*]")
        .appName("example")
        .config("spark.cassandra.connection.host","10.0.0.1") \
        .config("spark.cassandra.connection.port","9042") \
        .config("spark.cassandra.auth.username","pc1") \
        .config("spark.cassandra.auth.password","1234") \
        .config("spark.jars","/opt/spark/jars/spark-cassandra-connector_2.12-3.4.0.jar") \
        .getOrCreate()

同时需确认Connector版本与你的Spark(3.4.x)、Scala(2.12)版本完全匹配,版本不兼容会导致读写异常。

2. 规范写入Cassandra的语法

Spark Cassandra Connector写入时必须指定目标Keyspace和表名,仅调用df.write.save()会因缺少目标信息报错。正确的写入语法示例:

df.write \
    .format("org.apache.spark.sql.cassandra") \
    .option("keyspace", "your_target_keyspace") \
    .option("table", "your_target_table") \
    .mode("append")  # 根据需求选择append/overwrite/ignore等模式
    .save()

注意:必须确保DataFrame包含Cassandra表的所有主键列,且字段类型与表结构完全匹配(如Cassandra的int对应Spark的IntegerType,text对应StringType)。

3. 验证Cassandra权限与网络连通性

  • 权限检查:写入需要额外的WRITE权限,登录Cassandra执行以下命令确认:
    GRANT WRITE ON KEYSPACE your_target_keyspace TO pc1;
    
  • 网络检查:在Airflow所在的Docker容器内,执行telnet 10.0.0.1 9042验证9042端口是否连通。若不通,需检查Windows防火墙、Cassandra服务器安全组是否开放该端口,或Docker容器的网络模式是否能访问远程Cassandra。

4. 检查Airflow SparkSubmitOperator的配置

若使用SparkSubmitOperator提交任务,需确保在operator中传递正确的Spark配置,避免容器环境变量缺失导致配置失效:

from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

spark_write_task = SparkSubmitOperator(
    task_id="spark_write_to_cassandra",
    application="/path/to/your/spark_job.py",
    conn_id="spark_default",
    conf={
        "spark.cassandra.connection.host": "10.0.0.1",
        "spark.cassandra.connection.port": "9042",
        "spark.cassandra.auth.username": "pc1",
        "spark.cassandra.auth.password": "1234",
        "spark.jars": "/opt/spark/jars/spark-cassandra-connector_2.12-3.4.0.jar"
    }
)

内容的提问来源于stack exchange,提问作者zehra özdemir

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:02:22