【发布时间】:2018-02-12 14:17:09
【问题描述】:
我正在尝试使用来自 Datastax(v2.0.2、Spark v2.0.0)的 Spark-Cassandra 连接器:
val df = sparkSession.sparkContext.cassandraTable[MyRec](keyspace, tableName).toDF()
df.write.format("orc").save(hdfsLocation)
它看起来很简单并且工作了一段时间,但我开始遇到这样的异常:
Caused by: com.datastax.driver.core.exceptions.ReadFailureException:
Cassandra failure during read query at consistency LOCAL_ONE (1
responses were required but only 0 replica responded, 1 failed)
...
at com.datastax.spark.connector.rdd.CassandraTableScanRDD.com$datastax$
spark$connector$rdd$CassandraTableScanRDD$$fetchTokenRange(
CassandraTableScanRDD.scala:342)
增加spark.cassandra.read.timeout_ms和spark.cassandra.connection.timeout_ms和
减少spark.cassandra.input.fetch.size_in_rows 没有帮助。还玩了读一致性级别。
我在桌子上做了一个大的压实,但没有帮助。
因为这是一个产品。 DB 我无法调整服务器端参数,例如
tombstone_failure_threshold 建议here。
将完整表从 Cassandra (v3.7.0) 加载到 HDFS (Hive) 的最有效方法是什么?
【问题讨论】:
-
我认为这里的问题是在 Cassandra 方面,而不是在 Spark 上,也许这就是你所面临的:groups.google.com/a/lists.datastax.com/forum/#!topic/…
-
感谢链接。我同意这是 Cassandra 问题,很可能是因为墓碑。有没有办法仍然进行完全转储并避免这样的问题?使用 CqlInputFormat 的 MR 作业会更高效吗?
-
您知道 C* 和 Hive 的供应商吗? (Apache/HDP/CDH)
-
@saitejalakkimsetty Apache Cassandra 和 HDP 2.5.3
-
@Bruckwald,就像@RussS 建议的那样,这是由于峰值负载导致的临时可用性问题,这意味着查询/作业不能很好地扩展。在生产环境中,涉及流式传输/查询整个表的操作是一种反模式 (C*)。有很多方法可以解决这个问题。您可以通过将
CQL查询作为 `select * from test_table where token(PK) > (sometoken);` 来分步转储 C* 表,您可以处理读取超时和重试。另一种可扩展的解决方案是使用 Kafka 创建一个流式处理,因为它已包含在 HDP 中。
标签: scala hadoop apache-spark cassandra spark-dataframe