【发布时间】: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