【问题标题】:Avoid the use of Java data structures in Apache Spark to avoid copying the data避免在 Apache Spark 中使用 Java 数据结构以避免复制数据
【发布时间】:2016-06-13 04:16:58
【问题描述】:

我有一个 MySQL 数据库,其中包含大约 1 亿条记录(~25GB,~5 列)的单个表。使用 Apache Spark,我通过 JDBC 连接器提取这些数据并将其存储在 DataFrame 中。 从这里开始,我对数据进行了一些预处理(例如替换 NULL 值),所以我绝对需要遍历每条记录。 然后我想执行降维和特征选择(例如使用 PCA),执行聚类(例如 K-Means),然后在新数据上测试模型。

我已经在 Spark 的 Java API 中实现了这一点,但它太慢了(出于我的目的),因为我将数据从 DataFrame 复制到 java.util.Vector 和 java.util.List (到能够遍历所有记录并进行预处理),然后返回一个 DataFrame(因为 Spark 中的 PCA 需要一个 DataFrame 作为输入)。

我尝试将信息从数据库中提取到 org.apache.spark.sql.Column 中,但找不到对其进行迭代的方法。 我还尝试通过使用 org.apache.spark.mllib.linalg.{DenseVector, SparseVector} 来避免使用 Java 数据结构(例如 List 和 Vector),但也无法使其正常工作。 最后,我还考虑使用 JavaRDD(通过从 DataFrame 和自定义模式创建它),但无法完全解决。

经过冗长的描述,我的问题是:有没有一种方法可以完成第一段中提到的所有步骤,而无需将所有数据复制到 Java 数据结构中? 也许我尝试的其中一个选项实际上可以工作,但我似乎无法找到如何工作,因为关于 Spark 的文档和文献有点稀缺。

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-dataframe


    【解决方案1】:

    从您问题的措辞看来,Spark 处理的各个阶段似乎有些混乱。

    首先,我们通过指定输入和转换来告诉 Spark 要做什么。在这一点上,唯一已知的是 (a) 不同处理阶段的分区数量和 (b) 数据的模式。 org.apache.spark.sql.Column 在此阶段用于标识与列关联的元数据。但是,它不包含任何数据。事实上,现阶段根本没有数据。

    其次,我们告诉 Spark 对数据帧/数据集执行操作。这就是处理的开始。输入被读取并流经各种转换并进入最终操作操作,无论是collectsave 还是其他。

    所以,这就解释了为什么您不能“从数据库中提取信息到”Column

    至于您问题的核心,如果不查看您的代码并且确切地知道您要完成的工作是什么,就很难发表评论,但可以肯定地说,在类型之间进行大量迁移是一个坏主意。

    这里有几个问题可能会帮助您获得更好的结果:

    • 为什么不能通过直接在Row 实例上操作来执行所需的数据转换?

    • 将一些转换代码包装到 UDF 或 UDAF 中是否方便?

    希望这会有所帮助。

    【讨论】:

    • Sim,你写了一些非常有用的概念性的东西,这些对我很有用。同时。我设法通过在 Dataframes 上使用 UDF 来解决我的问题。此外,Spark 2.0.0 Preview 提供了许多附加功能,因此将我的代码移植到 2.0.0 可以更轻松地专门使用数据集(与将数据复制到 java Vector/List 相比)
    • 是的,如果你使用相对简单的数据结构,2.0.0 中的Datasets 应该可以正常工作。
    猜你喜欢
    • 1970-01-01
    • 2015-06-24
    • 2012-04-02
    • 2019-06-07
    • 2021-11-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-29
    相关资源
    最近更新 更多