Spark Cassandra Connector写入超时错误排查求助
背景信息
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

