【发布时间】:2020-01-18 13:08:32
【问题描述】:
我正在使用 spark-sql 2.4.1、spark-cassandra-connector_2.11-2.4.1.jar 和 java8。我有一种情况,出于审计目的,我需要计算 C* 表的表行数。 我的 C* 表中有大约 20 亿条记录。
为了计算行数,我尝试了两种方法,如下所示。
public static Long getColumnFamilyCountJavaApi(SparkSession spark,String keyspace, String columnFamilyName) throws IOException{
JavaSparkContext sc = new JavaSparkContext(spark.sparkContext());
return javaFunctions(sc).cassandraTable(keyspace, columnFamilyName).cassandraCount();
}
public static Long getColumnFamilyCount(SparkSession spark,String keyspace, String columnFamilyName) throws IOException{
return spark
.read()
.format("org.apache.spark.sql.cassandra")
.option("table", columnFamilyName)
.option("keyspace",keyspace )
.load().count();
}
但两种方式都会导致相同的错误。
Caused by: com.datastax.driver.core.exceptions.ReadFailureException: Cassandra failure during read query at consistency LOCAL_QUORUM (2 responses were required but only 0 replica responded, 2 failed)
at com.datastax.driver.core.exceptions.ReadFailureException.copy(ReadFailureException.java:85)
com.datastax.driver.core.DefaultResultSetFuture.getUninterruptibly(DefaultResultSetFuture.java:245)
at com.datastax.spark.connector.cql.DefaultScanner.scan(Scanner.scala:34)
at com.datastax.spark.connector.rdd.CassandraTableScanRDD.com$datastax$spark$connector$rdd$CassandraTableScanRDD$$fetchTokenRange(CassandraTableScanRDD.scala:342)
如何处理这种情况?
【问题讨论】:
标签: apache-spark cassandra apache-spark-sql datastax-enterprise datastax-java-driver