【问题标题】:How can I customize column-mapping when save RDD to Cassandra?将RDD保存到Cassandra时如何自定义列映射?
【发布时间】:2015-08-19 16:00:27
【问题描述】:

我正在使用 Java 编写 Spark 应用程序。 如果我有一个自定义元组,假设类“Person”。

Class Person { 
  public String name1; 
  public String name2; 
  public String name3; 
} 

我有一个

JavaRDD<Person> rdd;

现在我想将它保存到 Cassandra。

假设我在 Cassandra 中有一个名为“people”的表,其中包含“name1”、“name2”和“name3”、“name4”、...、“name10”三列。 根据教程,默认的列映射使用以下代码:

javaFunctions(rdd).writerBuilder("test", "person", mapToRow(Person.class)).saveToCassandra(); 

这将使用默认的列映射,例如:

Person.name1  --> "name1"    
Person.name2  --> "name2"     
Person.name3  --> "name3" 

但是我想自定义列映射,新的映射是这样的:

Person.name1  --> "name3"       
Person.name2  --> "name2"  
Person.name3  --> "name1" 

甚至我想丢弃 Person.name2

Person.name1  --> "name3"
Person.name3  --> "name1"

无论如何,我想知道是否有办法覆盖或替换默认的 RowWriter?
我应该如何修改列映射?
我找不到任何关于 Java 中自定义列映射的好材料。

【问题讨论】:

    标签: apache-spark spark-cassandra-connector


    【解决方案1】:

    请找到saveTOCassandra的签名

    def saveToCassandra(keyspaceName: String, 
                        tableName: String, columns: 
                        ColumnSelector = AllColumns, 
                         writeConf: WriteConf = WriteConf.fromSparkConf(sparkContext.getConf)) 
    

    解释:

    @param table 用于创建新表的表定义

    @param columns 选择要保存数据的列。 仅使用唯一的列名,并且您必须至少选择所有主键列。所有其他字段都被丢弃。未选择的属性/列名称保持不变。

    如果我正确理解您的需求,您可以使用参数“column”来实现您的结果。

    【讨论】:

    • 请采纳答案
    • 我的 cassandra 表列小写如下 CREATE TABLE model_family_by_id( model_family_id int PRIMARY KEY, model_family text, create_date date, last_update_date date, model_family_name text );
    • 我的数据框架构是这样的根 |-- MODEL_FAMILY_ID: decimal(38,10) (nullable = true) |-- MODEL_FAMILY: string (nullable = true) |-- CREATE_DATE: timestamp (nullable = true) |-- LAST_UPDATE_DATE: 时间戳 (nullable = true) |-- MODEL_FAMILY_NAME: string (nullable = true)
    • 因此,在线程“main”中插入 tabException 时 java.util.NoSuchElementException:在表 sample_cbd.model_family_by_id 中找不到列:COM.datastax.spark.connector 的 MODEL_FAMILY_ID、MODEL_FAMILY、CREATE_DATE、LAST_UPDATE_DATE、MODEL_FAMILY_NAME .SomeColumns.selectFrom(ColumnSelector.scala:44) 我得到错误
    • 我应该如何处理?
    猜你喜欢
    • 2016-01-30
    • 2017-04-10
    • 2015-03-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-09
    • 1970-01-01
    相关资源
    最近更新 更多