【问题标题】:Add Fields to Csv with Spark使用 Spark 将字段添加到 CSV
【发布时间】:2018-08-02 09:58:22
【问题描述】:

所以,我有一个 CSV,其中包含空间(latitudelongitude)和时间(timestamp)数据。

为了对我们有用,我们将空间信息转换为“geohash”,将时间信息转换为“timehash”。

问题是,如何使用 spark 将 geohashtimehash 添加为 CSV 中每一行的字段(因为数据约为 200 GB)?

我们尝试使用 JavaPairRDD 和它的函数 mapTopair ,但问题仍然在于如何转换回 JavaRdd 然后转换为 CSV?所以我认为这是一个糟糕的解决方案,我要求一种简单的方法。

问题更新:

在@Alvaro 的帮助下,我创建了这个 java 类:

public class Hash {
public static SparkConf Spark_Config;
public static JavaSparkContext Spark_Context;

UDF2 geohashConverter = new UDF2<Long, Long, String>() {
    
    public String call(Long latitude, Long longitude) throws Exception {
        // convert here
        return "calculate_hash";
    }
};

UDF1 timehashConverter = new UDF1<Long, String>() {
    
    public String call(Long timestamp) throws Exception {
        // convert here
        return "calculate_hash";
    }
};
public Hash(String path) {
    SparkSession spark = SparkSession
            .builder()
            .appName("Java Spark SQL Example")
            .config("spark.master", "local")
            .getOrCreate();
    
    spark.udf().register("geohashConverter", geohashConverter, DataTypes.StringType);
    spark.udf().register("timehashConverter", timehashConverter, DataTypes.StringType);
    
Dataset df=spark.read().csv(path)
    .withColumn("geohash", callUDF("geohashConverter", col("_c6"), col("_c7")))
    .withColumn("timehash", callUDF("timehashConverter", col("_c1")))
.write().csv("C:/Users/Ahmed/Desktop/preprocess2");

 }

public static void main(String[] args) {
    String path = "C:/Users/Ahmed/Desktop/cabs_trajectories/cabs_trajectories/green/2013";
    Hash h = new Hash(path);
}
}

然后我得到序列化问题,当我删除write().csv()时它消失了

【问题讨论】:

  • 您能否分享一些 sn-p 以便我们了解您当前如何加载 CSV?
  • 你需要知道我是如何在 spark 中加载 csv 的??!!
  • 我只是想看看您是使用 Datasets Spark API 还是只是创建和 RDD 逐行读取文件。

标签: java apache-spark apache-spark-sql rdd


【解决方案1】:

最有效的方法之一是使用数据集 API 加载 CSV 并使用用户定义的函数来转换您指定的列。这样,您的数据将始终保持结构,而不必处理元组。

首先,您创建用户定义函数:geohashConverter,它采用两个值(latitudelongitude)和 timehashConverter,它只采用时间戳。

UDF2 geohashConverter = new UDF2<Long, Long, String>() {
    @Override
    public String call(Long latitude, Long longitude) throws Exception {
        // convert here
        return "calculate_hash";
    }
};

UDF1 timehashConverter = new UDF1<Long, String>() {
    @Override
    public String call(Long timestamp) throws Exception {
        // convert here
        return "calculate_hash";
    }
};

创建后,您必须注册它们:

spark.udf().register("geohashConverter", geohashConverter, DataTypes.StringType);
spark.udf().register("timehashConverter", timehashConverter, DataTypes.StringType);

最后,只需读取您的 CSV 文件,并通过调用 withColumn 应用用户定义函数。它将根据您使用callUDF 调用的用户定义函数创建一个新列。 callUDF 总是收到一个字符串,其中包含您要调用的已注册 UDF 的名称以及一个或多个其值将传递给 UDF 的列。

最后,只需调用 write().csv("path") 即可保存您的数据集

import static org.apache.spark.sql.functions.col;
import static org.apache.spark.sql.functions.callUDF;


spark.read().csv("/source/path")
        .withColumn("geohash", callUDF("geohashConverter", col("latitude"), col("longitude")))
        .withColumn("timehash", callUDF("timehashConverter", col("timestamp")))
.write().csv("/path/to/save");

希望对您有所帮助!

更新

如果您发布导致问题的代码会非常有帮助,因为异常几乎没有说明代码的哪些部分不可序列化。

无论如何,根据我对 Spark 的个人经验,我认为问题在于您用于计算哈希的对象。请记住,该对象必须通过集群进行分发。如果这个对象不能被序列化,它会抛出一个Task not serializable Exception。您有两种解决方法:

  • 在用于计算哈希的类中实现Serializable 接口。

  • 创建一个生成哈希的静态方法并从 UDF 调用此方法。

更新 2

然后我得到序列化问题,当我删除时它消失了 写().csv()

这是预期的行为。当您删除 write().csv() 时,您什么也没有执行。你应该知道 Spark 是如何工作的。在这段代码中,csv() 之前调用的所有方法都是转换。在 Spark 中,直到调用 csv()show()count() 之类的操作后才会执行转换。

问题是您在不可序列化的类中创建和执行 Spark 作业(甚至在构造函数中更糟糕!!!??)

以静态方法创建 Spark 作业可以解决问题。请记住,您的 Spark 代码必须通过集群分发,因此它必须是可序列化的。它对我有用,必须对你有用:

public class Hash {
    public static void main(String[] args) {
        String path = "in/prueba.csv";

        UDF2 geohashConverter = new UDF2<Long, Long, String>() {

            public String call(Long latitude, Long longitude) throws Exception {
                // convert here
                return "calculate_hash";
            }
        };

        UDF1 timehashConverter = new UDF1<Long, String>() {

            public String call(Long timestamp) throws Exception {
                // convert here
                return "calculate_hash";
            }
        };

        SparkSession spark = SparkSession
                .builder()
                .appName("Java Spark SQL Example")
                .config("spark.master", "local")
                .getOrCreate();

        spark.udf().register("geohashConverter", geohashConverter, DataTypes.StringType);
        spark.udf().register("timehashConverter", timehashConverter, DataTypes.StringType);

        spark
                .read()
                .format("com.databricks.spark.csv")
                .option("header", "true")
                .load(path)
                .withColumn("geohash", callUDF("geohashConverter", col("_c6"), col("_c7")))
                .withColumn("timehash", callUDF("timehashConverter", col("_c1")))
                .write().csv("resultados");
    }
}

【讨论】:

  • 感谢您的帮助,但我的 java 编译器无法识别 'col()',是否需要特定的导入。
  • 试试看:import static org.apache.spark.sql.functions.col;
  • 另一个问题 callUDF 的属性不清楚我在the doc 中找不到解释
  • callUDF 总是收到一个字符串,其中包含您要调用的已注册 UDF 的名称以及一个或多个其值将传递给 UDF 的列。
  • 我遇到序列化问题我会更新我的问题
猜你喜欢
  • 1970-01-01
  • 2015-07-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-09-13
  • 2019-12-08
  • 2019-05-09
  • 2017-05-15
相关资源
最近更新 更多