spark-cassandra-connector调用repartitionByCassandraReplica返回空RDD问题
问题原因
最常见原因:缺失本地DC配置
如果你的Cassandra集群使用NetworkTopologyStrategy作为keyspace的复制策略,而Spark配置中没有指定spark.cassandra.connection.localDC参数,repartitionByCassandraReplica方法无法匹配到对应副本的位置,就会生成分区数为0的空RDD。
次常见原因:3.0.0版本Connector参数bug
Spark Cassandra Connector 3.0.0版本存在已知问题:调用repartitionByCassandraReplica时显式传入自定义分区数参数,会导致令牌环分区计算错误,过滤掉所有有效数据。
其他可能原因
- 输入的
experimentlist中存在空值,或者experimentid字符串的大小写和Cassandra表中存储的值不匹配(Cassandra varchar类型大小写敏感) ExperimentForm类的Bean序列化规则和Connector的映射逻辑不匹配,导致分区键值提取失败
解决方案
- 添加本地DC配置
在Spark会话初始化时添加如下配置,DC名称可以通过在Cassandra节点执行nodetool status命令查询:
spark.conf().set("spark.cassandra.connection.localDC", "你的Cassandra集群DC名称");
- 去掉显式分区数参数
修改repartitionByCassandraReplica的调用方式,使用默认分区数的方法重载:
JavaRDD<ExperimentForm> predf = CassandraJavaUtil.javaFunctions(dfexplistoriginal.toJavaRDD()) .repartitionByCassandraReplica("mdb","experiment", CassandraJavaUtil.someColumns("experimentid"), CassandraJavaUtil.mapToRow(ExperimentForm.class));
- 校验输入数据
提前过滤experimentlist中的空值,确认输入的experimentid大小写和Cassandra表中的存储值完全一致。
优化建议
你当前的实现逻辑不需要将RDD转Dataset后再和全量Cassandra表做关联,直接使用joinWithCassandraTable方法可以避免全表扫描和额外shuffle,性能提升明显:
JavaPairRDD<ExperimentForm, Row> joinedRDD = CassandraJavaUtil.javaFunctions(dfexplistoriginal.toJavaRDD()) .joinWithCassandraTable("mdb", "experiment", CassandraJavaUtil.someColumns("experimentid"), CassandraJavaUtil.someColumns("experimentid", "description", "intensity"), CassandraJavaUtil.mapToRow(ExperimentForm.class), Row.class); // 可直接将joinedRDD转为需要的Dataset使用
内容的提问来源于stack exchange,提问作者ktzan
相关产品推荐
相关产品推荐

