【问题标题】:Using map function in Apache Spark for huge operation在 Apache Spark 中使用 map 函数进行大型操作
【发布时间】:2017-07-30 11:15:35
【问题描述】:

我们需要像 jaccard 一样在 spark 中的大量数据集上计算距离矩阵。 面临几个问题。请帮我们指路。

问题 1

    import info.debatty.java.stringsimilarity.Jaccard;

    //sample Data set creation
    List<Row> data = Arrays.asList(
                RowFactory.create("Hi I heard about Spark", "Hi I Know about Spark"),
                RowFactory.create("I wish Java could use case classes","I wish C# could use case classes"),
                RowFactory.create("Logistic,regression,models,are,neat","Logistic,regression,models,are,neat"));

    StructType schema = new StructType(new StructField[] {new StructField("label", DataTypes.StringType, false,Metadata.empty()),
                new StructField("sentence", DataTypes.StringType, false,Metadata.empty()) });
                Dataset<Row> sentenceDataFrame = spark.createDataFrame(data, schema);

                // Distance matrix object creation
                Jaccard jaccard=new Jaccard();

                //Working on each of the member element of dataset and applying distance matrix.
                Dataset<String> sentenceDataFrame1 =sentenceDataFrame.map(
                        (MapFunction<Row, String>) row -> "Name: " + jaccard.similarity(row.getString(0),row.getString(1)),Encoders.STRING()
                );
                sentenceDataFrame1.show();

没有编译时错误。但是出现运行时异常,例如:

org.apache.spark.SparkException: Task not serializable

问题 2
此外,我们需要找到哪对得分最高,我们需要声明一些变量。此外,我们还需要执行其他计算,我们面临很多困难。
即使我尝试在 MapBlock 中声明一个像计数器这样的简单变量,我们也无法捕获增量值。如果我们在 Map 块之外声明,我们会得到很多编译时错误。

    int counter=0;
        Dataset<String> sentenceDataFrame1 =sentenceDataFrame.map(
                (MapFunction<Row,  String>) row -> {
                    System.out.println("Name: " + row.getString(1));
                    //int counter = 0;
                    counter++;
                    System.out.println("Counter: " + counter);
                    return counter+"";

                },Encoders.STRING()

        );

请给我们指路。 谢谢。

【问题讨论】:

    标签: apache-spark dataset similarity apache-spark-2.0 apache-spark-dataset


    【解决方案1】:

    Jaccard jaccard=new Jaccard();

    这个类可以序列化吗?

    在 spark 中,您在 Transformations 中编写的所有代码都在驱动程序上实例化并序列化并发送到执行程序。

    你已经使用了 lambda 函数:

    1. lambda 内部外部类使用的所有类都需要可序列化。

    2. 如果甚至使用 lambda 内部的外部类中的方法,它期望外部类是可序列化的。

    详细了解请参考:

    http://bytepadding.com/big-data/spark/spark-code-analysis/

    http://bytepadding.com/big-data/spark/understanding-spark-serialization/

    第 2 部分:

    1. 尝试在 spark 中查找笛卡尔积 N CROSS N。
    2. 尝试寻找更聪明的算法来找到配对。

    对该问题的更多输入将有助于提供更好的答案。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-02-17
      • 2023-04-03
      • 1970-01-01
      • 2014-01-29
      • 1970-01-01
      • 2016-11-19
      • 1970-01-01
      相关资源
      最近更新 更多