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

使用Spark writeStream无法向Cassandra插入数据的问题求助

流处理管道Cassandra数据插入问题排查

问题描述

搭建流处理管道:Airflow调用API -> Kafka处理 -> Spark写入Cassandra。目前已成功创建keyspace和表,但无数据插入,Airflow的DAG已确认将数据正确发送到控制服务器。

相关代码片段:

if __name__ == "__main__":
    # create spark connection
    spark_conn = create_spark_connection()

    if spark_conn is not None:
        # connect to kafka with spark connection
        spark_df = connect_to_kafka(spark_conn)
        selection_df = create_selection_df_from_kafka(spark_df)
        session = create_cassandra_connection()

        if session is not None:
            create_keyspace(session)
            create_table(session)
            """"""""
            logging.info("Streaming is being started...")
            
            streaming_query = (selection_df.writeStream.format("org.apache.spark.sql.cassandra")
                               .option('checkpointLocation', '/tmp/checkpoint')
                               .option('keyspace', 'spark_streams')
                               .option('table', 'created_users')
                               .start())

            streaming_query.awaitTermination()
            """"""""

排查建议

  • 修复代码注释问题:你的流处理核心逻辑被多引号""""""注释掉了,这部分代码根本没有执行,直接删除多余的引号让逻辑运行起来。
  • 验证Spark是否从Kafka接收数据:在selection_df创建后添加临时控制台输出,确认数据是否到达Spark:
    selection_df.writeStream.outputMode("append").format("console").start().awaitTermination()
    
    同时检查connect_to_kafka中的Kafka主题、消费组、偏移量配置是否正确。
  • 检查Spark-Cassandra依赖与配置:
    • 确认Spark作业依赖了对应版本的spark-cassandra-connector和spark-sql-kafka-0-10包,版本需与Spark、Cassandra兼容。
    • 检查checkpointLocation指定的/tmp/checkpoint路径是否有写入权限,无权限会导致流任务静默失败。
  • 查看日志定位错误:
    • 开启Spark DEBUG日志,查看流任务启动、数据写入阶段的日志,排查是否有Cassandra连接失败、数据格式不匹配等报错。
    • 检查Cassandra系统日志,确认是否有写入请求到达但被拒绝(如权限不足、一致性级别不满足)。
  • 验证数据结构匹配:
    • 对比selection_df的Schema与Cassandracreated_users表的字段名、数据类型,确保完全对应(如Cassandra的text对应Spark的StringType)。
    • 确认selection_df包含Cassandra表的所有主键字段,且字段值不为空。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 18:34:58