【问题标题】:Getting null pointer exception when trying to add a column in Spark Dataset in Java尝试在 Java 中的 Spark 数据集中添加列时出现空指针异常
【发布时间】: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


    【解决方案1】:

    虽然使用 df.withColumn 需要第一个参数作为列名,第二个参数作为该列的值。 如果您想添加名称为“cluster”且值来自某个 json 值的新列,则可以使用“lit”函数作为 lit(cluster_val),其中 cluster_val 保存值。

    您必须导入“org.apache.spark.sql.functions._”才能使用 lit 函数。

    希望对你有帮助。

    【讨论】:

    • 感谢@Ramdev Sharma。如果我使用df.withColumn("cluster", functions.lit(cluster_val));,我仍然会得到空指针异常,即使我有cluster_val 的有效值
    • 你能在评论 df.withColumn 行后检查 df.show 是否工作吗?
    • 你猜对了。 df.show() 在我的代码中的 df.foreach((ForeachFunction&lt;Row&gt;) row -&gt;...{..} 代码块之后不起作用。在此行之前, df.show() 有效。
    • 所以我尝试将数据集转换移到循环之外,但是这样做时你得到了新列的值(检查我有问题的编辑),我得到了空值。为什么?
    • 您可以为输入和预期输出添加示例数据吗?
    猜你喜欢
    • 2019-10-16
    • 2023-03-18
    • 2014-03-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多