【发布时间】:2016-02-24 00:59:06
【问题描述】:
我正在尝试将 Spark RDD 中的 cassandra 行列映射到我可以交互以在 spark 中进行操作但似乎无法将它们放入变量的变量。我有以下代码:
JavaRDD<MeasuredValue> rdd = javaFunctions(sc).cassandraTable("model", "reports", mapRowTo (MeasuredValue.class))
.select("start_frequency","bandwidth", "power");
JavaRDD<Value> valueRdd = rdd.flatMap(row-> {
double start_frequency = row.getStartFrequency();
float power = row.getPower();
double bandwidth = row.getBandwidth();
List<Value> list = new ArrayList<Value>();
// Create Channel Power Buckets
for(channel = 1.6000E8; channel <= channel_end; ){
if( (channel >= start_frequency) && (channel <= (start_frequency + bandwidth)) ) {
list.add(new Value(channel, power));
} // end if
channel+=increment;
} // end for
})
我的课程如下所示:
public class Value implements Serializable {
public Value(Double channel, Float power) {
this.channel = channel;
this.power = power;
}
Double channel;
Float power;
public void setChannel(Double channel) {
this.channel = channel;
}
public void setPower(Float power) {
this.power = power;
}
public Double getChannel() {
return channel;
}
public Float getPower() {
return power;
}
@Override
public String toString() {
return "[" +channel +","+power+"]";
}
}
public static class MeasuredValue implements Serializable {
public MeasuredValue() { }
private double start_frequency;
public double getStart_frequency() { return start_frequency; }
public void setStart_frequency(double start_frequency) { this.start_frequency = start_frequency; }
private double bandwidth ;
public double getBandwidth() { return bandwidth; }
public void setBandwidth(double bandwidth) { this.bandwidth = bandwidth; }
private float power;
public float getPower() { return power; }
public void setPower(float power) { this.power = power; }
}
我尝试使用 lambda 对行进行平面映射的尝试似乎是错误的,因为我收到以下错误:
AbstractJavaRDDlike 类中的方法 flatMap 无法应用 给定类型;必需:找到 FlatMapFunction: (row)->{d[...];}} 原因:无法推断类型变量 U(参数 不匹配; lambda 表达式中的错误返回类型缺少返回值)
我在“创建通道电源桶”循环中遇到了关于
的错误"从 lambda 表达式引用的局部变量必须是 final 或实际上是最终的”
如果我可以使用 DataFrame 来做到这一点,我会对查看代码来促进这一点感兴趣。
【问题讨论】:
-
我应该使用DataFrame而不是RDD吗?
-
第二条错误消息表明 lambda 中使用的某些变量未声明为 final -
increment和channel_end变量是什么?他们是final吗? -
它们的定义如下:
// Define Variable double channel,channel_end,start_frequency, increment, bandwidth; float power; long time_key; // Initialize Variables channel_end = 1.6159E8; increment = 5000; -
好吧,就是这样(或其中的一部分)——它们必须是最终的,例如
final double channel_end = 1.6159E8; -
主要问题是能够将行列值映射到变量。我可以从火花中操纵。
标签: java apache-spark cassandra datastax-java-driver