【发布时间】:2018-10-08 18:53:31
【问题描述】:
我正在尝试遍历 Java 中的 Dataset 行,然后访问特定列以查找其存储为 JSON 文件中的键的值并获取其值。找到的值需要存储为该行中所有行的新列值。
我看到我从 JSON 文件获得的cluster_val 不是 NULL,但是当我尝试将其添加为列时,我得到了Exception in thread "main" org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1.0 failed 1 times, most recent failure: Lost task 0.0 in stage 1.0 (TID 1, localhost, executor driver): java.lang.NullPointerException
到目前为止,我有这个:
Dataset<Row> df = spark.read().format("csv").load(path);
df.foreach((ForeachFunction<Row>) row ->
{
String df_col_val = (String) row.get(6);
System.out.println(row.get(6));
if(df_col_val.length() > 5){
df_col_val = df_col_val.substring(0, df_col_val.length() - 5 + 1); //NOT NULL
}
System.out.println(df_col_val);
String cluster_val = (String) jo.get(df_col_val); //NOT NULL
System.out.println(cluster_val);
df.withColumn("cluster", df.col(cluster_val)); // NULL POINTER EXCEPTION. WHY?
df.show();
});
所以我主要需要帮助逐行读取数据集并执行上述后续操作。 网上找不到太多参考资料。如果可能,请向我推荐正确的来源。另外,如果有速记方法,请告诉我。
所以我发现df.col(cluster_val) 正在抛出异常,因为没有现有的列。如何将列的字符串名称转换为传入withColumn()函数pf数据集所需的列类型
更新:
所以我尝试了以下方法,在这里我尝试使用 udf 获取新列的值,但如果像这样使用它则为空:
Dataset<Row> df = spark.read().format("csv").option("header", "true").load(path);
Object obj = new JSONParser().parse(new FileReader("path to json"));
JSONObject jo = (JSONObject) obj;
df.withColumn("cluster", functions.lit((String) jo.get(df.col(df_col_val)))));
df.show();
【问题讨论】:
标签: java loops apache-spark apache-spark-dataset