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

Spark Cassandra Connector写入超时错误排查求助

Spark写入Cassandra超时错误分析

背景信息

Cassandra表定义

CREATE TABLE my_keyspace.my_table (
    my_composite_pk_a bigint,
    my_composite_pk_b ascii,
    value blob,
    PRIMARY KEY ((my_composite_pk_a, my_composite_pk_b))
) WITH bloom_filter_fp_chance = 0.1
   AND gc_grace_seconds = 86400
   AND caching = {'keys': 'ALL', 'rows_per_partition': 'NONE'}
   AND compaction = {'class': 'org.apache.cassandra.db.compaction.LeveledCompactionStrategy', 'enabled': 'true'}
   AND compression = {'chunk_length_in_kb': '64', 'class': 'org.apache.cassandra.io.compress.LZ4Compressor'};

注:value字段为BLOB类型,单条数据约1MB

Spark写入逻辑

spark
  // 从parquet读取数据
  .read.parquet("...")
  // 跳过大数据分区,避免压垮Cassandra
  .withColumn("bytes_count",length(col("value")))
  .filter("bytes_count < 1000000") // 小于1MB
  // 投影字段
  .select("my_composite_pk_a", "my_composite_pk_b", "value")
  // 写入Cassandra
  .writeTo("cassandra.my_keyspace.my_table")
  .append()

Spark Cassandra Connector配置

spark.sql.catalog.cassandra.spark.cassandra.output.concurrent.writes=6 
spark.sql.catalog.cassandra.spark.cassandra.output.batch.size.rows=1 
spark.sql.catalog.cassandra.spark.cassandra.output.batch.grouping.key=none 
spark.sql.catalog.cassandra.spark.cassandra.output.throughputMBPerSec=6 
spark.sql.catalog.cassandra.spark.cassandra.connection.host=node1,node2 
spark.sql.catalog.cassandra.spark.cassandra.connection.port=9042 
spark.sql.catalog.cassandra.spark.cassandra.output.consistency.level=LOCAL_QUORUM 
spark.sql.catalog.cassandra.spark.cassandra.output.metrics=false 
spark.sql.catalog.cassandra.spark.cassandra.connection.timeoutMS=90000 
spark.sql.catalog.cassandra.spark.cassandra.query.retry.count=100 
spark.sql.catalog.cassandra=com.datastax.spark.connector.datasource.CassandraCatalog 
spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions 
spark.sql.catalog.cassandra.spark.cassandra.auth.username=USERNAME 
spark.sql.catalog.cassandra.spark.cassandra.auth.password=PASSWORD

Spark集群配置

--total-executor-cores 6 
--executor-cores 6 
--executor-memory 15G 
--driver-memory 6G 
--driver-cores 4 

注:当前仅使用1个Executor,配备6核

错误信息

...
24/03/06 10:07:24 WARN TaskSetManager: Lost task 886.0 in stage 3.0 (TID 2959, node1, executor 0): java.util.concurrent.TimeoutException: Futures timed out after [10 seconds]
...
Caused by: java.util.concurrent.TimeoutException: Cannot receive any reply from node4:40201 in 10 seconds
...

可能的原因

  • 写入超时参数不匹配
    虽然配置了spark.cassandra.connection.timeoutMS=90000(连接超时90秒),但Cassandra节点默认的write_request_timeout_in_ms是10000ms(10秒),而单条1MB的Blob写入操作耗时可能超过这个阈值。同时Spark Connector未显式设置spark.cassandra.output.writeTimeoutMS,会默认沿用Cassandra节点的超时配置,导致写入请求在10秒后被判定为超时。

  • 批量写入配置完全失效
    当前设置spark.cassandra.output.batch.size.rows=1和spark.cassandra.output.batch.grouping.key=none,意味着每条1MB的Blob数据都会单独发起一个写入请求,没有任何批量合并优化。这会产生海量的单点写入请求,Cassandra节点需要频繁处理单个请求的网络握手、磁盘IO,尤其是node4可能承担了较多目标分区的写入压力,最终因负载过高无法及时响应。

  • Cassandra节点资源瓶颈
    错误明确指向node4无响应,该节点大概率存在资源耗尽的情况:

    • 磁盘IO瓶颈:LeveledCompactionStrategy(LCS)在处理大Blob数据时,会频繁执行压缩和分区合并操作,占用大量磁盘IO带宽,导致写入请求无法及时落盘。
    • CPU负载过高:LZ4压缩、分区索引维护等操作占用过多CPU资源,节点无法及时处理新的写入请求。
    • 内部网络问题:node4的内部通信端口(40201)出现阻塞,或者与Spark Executor之间的网络存在延迟、丢包。
  • 并发与吞吐量配置的矛盾
    spark.cassandra.output.concurrent.writes=6和throughputMBPerSec=6的配置,理论上刚好匹配单条1MB数据的6并发写入,但实际场景中,磁盘IO延迟可能远超预期,导致单条写入耗时超过10秒,并发请求堆积在Cassandra节点,最终触发超时。

  • Spark任务执行效率不足
    仅使用1个Executor处理所有写入任务,即使配备6核,也可能因单节点网络连接资源耗尽、任务排队等问题,导致无法及时处理Cassandra的响应,间接加剧超时问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 19:35:13