【问题标题】:Spark with Scala: write null-like field value in Cassandra instead of TupleValueSpark 与 Scala:在 Cassandra 中而不是 TupleValue 中写入类似空的字段值
【发布时间】:2017-05-11 19:54:26
【问题描述】:

在我的一个收藏中,假设我有以下字段:

f: frozen<tuple<text, set<text>>

假设我想使用 Scala 脚本在该特定字段为空、null、不存在等的地方插入一个条目,在插入之前我将条目的字段映射如下:

sRow("fk") = null // or None, or maybe I simply don't specify the field at all

当尝试运行 spark 脚本(来自 Databricks,Spark 连接器版本 1.6)时,我收到以下错误:

org.apache.spark.SparkException: Job aborted due to stage failure: Task 6 in stage 133.0 failed 1 times, most recent failure: Lost task 6.0 in stage 133.0 (TID 447, localhost): com.datastax.spark.connector.types.TypeConversionException: Cannot convert object null to com.datastax.spark.connector.TupleValue.
    at com.datastax.spark.connector.types.TypeConverter$$anonfun$convert$1.apply(TypeConverter.scala:47)
    at com.datastax.spark.connector.types.TypeConverter$$anonfun$convert$1.apply(TypeConverter.scala:43)

当使用None 而不是null 时,我仍然得到一个错误,虽然是一个不同的错误:

org.apache.spark.SparkException: Job aborted due to stage failure: Task 2 in stage 143.0 failed 1 times, most recent failure: Lost task 2.0 in stage 143.0 (TID 474, localhost): java.lang.IllegalArgumentException: requirement failed: Expected 2 components, instead of 0
    at scala.Predef$.require(Predef.scala:233)
    at com.datastax.spark.connector.types.TupleType.newInstance(TupleType.scala:55)

我知道 Cassandra 没有确切的 null 概念,但我知道在将条目插入 Cassandra 时有一种方法可以将值排除在外,就像我在其他环境中所做的那样,例如为 Cassandra 使用 nodejs 驱动程序. 在插入预期的 TupleValue 或某些用户定义的类型时,如何强制使用类似 null 的值?

【问题讨论】:

    标签: scala apache-spark cassandra databricks


    【解决方案1】:

    使用现代版本的 Cassandra,您可以使用“未绑定”功能让它实际跳过空值。这可能最适合您的用例,因为编写 null 会隐式写入墓碑。

    Treating nulls as Unset

    //Setup original data (1, 1, 1) --> (6, 6, 6)
    sc.parallelize(1 to 6).map(x => (x, x, x)).saveToCassandra(ks, "tab1")
    
    val ignoreNullsWriteConf = WriteConf.fromSparkConf(sc.getConf).copy(ignoreNulls = true)
    //These writes will not delete because we are ignoring nulls
    val optRdd = sc.parallelize(1 to 6)
      .map(x => (x, None, None))
      .saveToCassandra(ks, "tab1", writeConf = ignoreNullsWriteConf)
    
    val results = sc.cassandraTable[(Int, Int, Int)](ks, "tab1").collect
    
    results
    /**
      (1, 1, 1),
      (2, 2, 2),
      (3, 3, 3),
      (4, 4, 4),
      (5, 5, 5),
      (6, 6, 6)
    **/
    

    还有更细粒度的控件 Full Docs

    【讨论】:

    • 如果您希望 null 被视为未设置
    • Treating nulls as Unset 的链接文档特别提到了 DataFrames。它也适用于 RDD 吗?
    • 文档说参数是Dataframes的方法,因为Dataframes不能使用其他方法。该参数和所有其他示例都适用于 RDDsl。
    猜你喜欢
    • 1970-01-01
    • 2019-04-09
    • 1970-01-01
    • 2018-03-11
    • 2019-11-12
    • 2015-04-14
    • 2018-12-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多