GCP Dataproc Serverless中Cassandra写入Bigtable异常求助
问题描述
- 任务目标:将Cassandra实例中带
LIMIT 1000的查询结果写入Bigtable,基于GCP Dataproc Python模板(Cassandra转BigQuery、GCS转Bigtable)开发 - 异常情况:任务运行超2小时未完成,Spark UI显示ShuffleMapStage有6372个缺失任务;Bigtable无任何数据写入;预期任务仅需数分钟完成
- 已排查:容器内
hbase-site.xml配置未发现明显错误
相关代码片段
spark = ( SparkSession.builder.appName("CassandraToBT") .config(constants.SQL_EXTENSION, constants.CASSANDRA_EXTENSION) .config(f"spark.sql.catalog.casscon", constants.CASSANDRA_CATALOG) .config( f"spark.sql.catalog.{catalog}.spark.cassandra.connection.host", input_host, ) .config("spark.cassandra.auth.username", input_username) .config("spark.cassandra.auth.password", input_password) .getOrCreate() ) print(spark.sql("SHOW NAMESPACES FROM casscon").show(), flush=True) print("***************************************************************", flush=True) print(f"Spark version:{spark.version}") spark.sql("set spark.sql.shuffle.partitions=4") # Read if not query: input_data = spark.read.table(f"{catalog}.{input_keyspace}.{input_table}") else: input_data = spark.sql(query) print("---------------------------------------------------------------", flush=True) # Write print(input_data.head(), flush=True) input_data.write.format(constants.FORMAT_HBASE).options(catalog=catalog).option( "hbase.spark.use.hbasecontext", "false" ).mode("overwrite").save()
问题排查与修复建议
1. Shuffle分区配置未生效
你使用spark.sql("set spark.sql.shuffle.partitions=4")设置分区,但这种会话级临时配置可能在读取Cassandra时未生效。Cassandra Spark连接器默认会根据数据量生成大量分片,即便查询带LIMIT,也可能先全表扫描再过滤,导致分片过多。
修复: 在SparkSession构建阶段直接配置分区和Cassandra读取分片大小:
spark = ( SparkSession.builder.appName("CassandraToBT") .config(constants.SQL_EXTENSION, constants.CASSANDRA_EXTENSION) .config(f"spark.sql.catalog.casscon", constants.CASSANDRA_CATALOG) .config(f"spark.sql.catalog.{catalog}.spark.cassandra.connection.host", input_host) .config("spark.cassandra.auth.username", input_username) .config("spark.cassandra.auth.password", input_password) .config("spark.sql.shuffle.partitions", "4") # 全局设置Shuffle分区 .config("spark.cassandra.input.split.size_in_mb", "64") # 限制Cassandra读取分片大小 .config("spark.cassandra.sql.pushdown.enable", "true") # 开启查询下推,让LIMIT生效在Cassandra端 .getOrCreate() )
2. Bigtable写入配置错误
当前代码复用Cassandra的catalog变量写入Bigtable,这是错误的——Bigtable的Spark连接器需要单独定义包含列族映射的catalog,且需确保关键配置正确。
修复:
- 定义Bigtable专属catalog(根据你的表结构调整列族和列映射):
bigtable_catalog = '''{ "table": {"namespace": "default", "name": "your-bigtable-table-name"}, "rowkey": "your-rowkey-column", "columns": { "column1": {"cf": "your-column-family", "col": "column1"}, "column2": {"cf": "your-column-family", "col": "column2"} } }'''
- 修改写入代码:
input_data.write.format(constants.FORMAT_HBASE) .option("catalog", bigtable_catalog) .option("hbase.spark.use.hbasecontext", "false") .mode("overwrite") .save()
- 确保
hbase-site.xml中包含以下配置,且Dataproc集群能加载到该文件(提交任务时可通过--files hbase-site.xml指定):
<property> <name>google.bigtable.project.id</name> <value>your-gcp-project-id</value> </property> <property> <name>google.bigtable.instance.id</name> <value>your-bigtable-instance-id</value> </property>
3. LIMIT查询未下推至Cassandra
即便SQL语句带LIMIT 1000,如果未开启查询下推,Spark会先全量读取Cassandra表数据,再在本地过滤,导致读取数据量远超预期,引发大量Shuffle任务。
修复:
- 开启
spark.cassandra.sql.pushdown.enable=true(已在第一条修复中包含) - 或者改用Cassandra数据源直接读取并限制条数:
input_data = spark.read.format("org.apache.spark.sql.cassandra") .options(table=input_table, keyspace=input_keyspace) .load() .limit(1000)
4. 任务缺失与数据未写入的深层原因
大量Shuffle任务缺失通常与资源不足、写入阻塞有关:
- 检查Dataproc集群Worker节点的CPU、内存配置,确保资源足够支撑任务
- 查看Spark Driver和Worker日志,定位任务失败的具体原因(如Bigtable连接超时、权限不足、列映射错误)
- 临时降低日志级别至DEBUG,排查数据读取和写入的中间流程
内容的提问来源于stack exchange,提问作者Adam C. Scott
相关产品推荐
相关产品推荐

