【问题标题】:Cassandra/Parquet union RDDCassandra/Parquet union RDD
【发布时间】:2015-05-27 10:55:52
【问题描述】:

我刚刚开始使用 spark-cassandra 连接器并遇到以下问题:我有一个数据集,部分位于 cassandra 中,部分位于 HDFS 中(完全相同的架构)。我想创建两个集合的单个 UnionRDD 并从那里继续。

到目前为止我的代码是这样的:

DataFrame df = sqlContext.parquetFile("foo.parquet");
JavaRDD cassandraRDD = (JavaRDD)javaFuntions(sc).cassandraTable("foo_ks","foo_table");
DataFrame cassandraDF = sqlContext.applySchema(cassandraRDD,df.schema());

我收到一个运行时错误,提示无法将 CassandraRow 强制转换为 spark.sql.Row,来自 applySchema 调用......这并不奇怪。什么是让它工作的正确方法? (我的最终目标是联合 df 和 cassandraDF)。

我正在尝试使用 Spark 1.3.1 和 cassandra-spark 的主分支构建。

【问题讨论】:

  • 如果遇到异常,则首先打印架构并与 cassandraRDD 的字段进行比较。

标签: cassandra apache-spark


【解决方案1】:

最简单的方法是编写一个地图函数,它将采用

  1. 卡桑德拉街
  2. 源架构对象
  3. 目标架构对象

这个地图功能将

  1. 使用源架构读取 cassandra 行(并处理问题,例如填充缺失的列、抑制存在某些数据质量问题的行等)
  2. 将 cassandra 架构转换为 spark sql 架构(这是一个静态映射黑白 cassandra 类型到 sql 类型)
  3. 返回带有目标架构的 SQL Row 对象

所以,你应该可以做喜欢的事

cDF = cRDD.map(c2r).createDataFrame() // map 会返回 row 所以这里不需要 applySchema

基本上,我建议使用单个函数来处理转换。一旦你从 cassandra 数据“创建”了一个 DF,你就可以与任何其他 DF 进行联合。

【讨论】:

  • 谢谢 ayan -- 我希望不必手动编写 c2r,因为我的行有 70 多个字段......反正已经在 cassandra 中输入了这些字段。我会将您的答案标记为已接受,因为我认为也没有更简单的方法...
  • 好吧,您可能不需要“手工”完成。您可以打开一个单独的 JDBC 连接并从 Cassandra 获取架构信息。然后你可以在 c2r 中使用它。这样,即使架构发生更改,您也无需更改代码。核心点吧,是的,我们要告诉spark关于schema。顺便说一句,新的 Cassandra 连接器已经推出,至少我听说过。你可以/应该看看....
猜你喜欢
  • 2016-02-18
  • 1970-01-01
  • 2019-11-12
  • 1970-01-01
  • 2016-01-30
  • 2020-06-25
  • 2016-10-14
  • 2017-04-10
  • 1970-01-01
相关资源
最近更新 更多