【问题标题】:How to specify BigDecimal scale and precision in schema when loading a Mongo collection as a Spark Dataset将 Mongo 集合加载为 Spark 数据集时如何在模式中指定 BigDecimal 比例和精度
【发布时间】:2020-08-03 09:58:30
【问题描述】:

我正在尝试使用 Scala Mongo 连接器将大型 Mongo 集合加载到 Apache Spark。

我正在使用以下版本:

libraryDependencies += "org.apache.spark" %% "spark-core" % "3.0.0" 
libraryDependencies += "org.apache.spark" %% "spark-sql" % "3.0.0" 
libraryDependencies += "org.mongodb.spark" %% "mongo-spark-connector" % "2.4.2"
scalaVersion := "2.12.12"
openjdk version "11.0.8" 2020-07-14

该集合包含大于1e13 的大整数十进制值。我想要获取的数据集是一个集合,其中包含一个名为Output 的相应案例类,定义为:

case class Output(time: Long, pubKeyId: Long, value: BigDecimal, outIndex: Long, outTxId: Long)

如果我使用 MongoSpark.load 而不指定案例类:

val ds = MongoSpark.load(sc, rc).toDS[Output]

然后 Mongo 通过随机抽样集合来推断模式。这导致value 属性的随机比例,并且value 溢出随机获得的比例的任何文档在结果数据集中都缺少value 属性。这显然是不希望的。

或者,根据documentation for the Mongo Spark connector,我可以通过将案例类指定为load 的类型参数来显式设置架构,例如:

val ds = MongoSpark.load[Output](sc, rc).toDS[Output]

但是,在 case-class 定义中,我只能将 value 的类型指定为 BigDecimal,这不允许我明确说明所需的比例和精度。生成的架构使用 (38,18) 的默认精度和小数位数,这并不总是需要的:

root
 |-- time: long (nullable = false)
 |-- pubKeyId: long (nullable = false)
 |-- value: decimal(38,18) (nullable = true)
 |-- outIndex: long (nullable = false)
 |-- outTxId: long (nullable = false)

这与 Spark SQL API 不同,后者允许使用 DecimalType 显式指定比例和精度,例如:

val mySchema = StructType(StructField("value", DecimalType(30, 0)) :: Nil)

在将 Mongo 集合加载到 Apache Spark 时,如何为架构中的大十进制类型请求特定的比例和精度,类似于上面的代码?

【问题讨论】:

    标签: mongodb apache-spark-sql schema precision bigdecimal


    【解决方案1】:

    据我所知,根据thisthis,Decimal128 中的尾数和指数是固定大小的。除非您能找到相反的证据,否则 MongoDB 允许为其小数指定小数位数和精度是没有意义的。

    我的理解是关系数据库会根据规模和精度使用不同的浮点类型(例如 32 位与 64 位浮点数),但在 MongoDB 中,数据库会保留它给定的类型,所以如果你想要更短的浮点数,你需要让您的应用程序发送它而不是十进制类型。

    【讨论】:

    • 问题不在于值如何在 Mongo 内部存储,而在于它在 Spark 应用程序中的表示方式。 Spark SQL Schema 指定 Spark 如何表示数据,而不是数据在 Mongo 中的存储方式。在 Scala Spark 应用程序中,当从 Mongo 加载数据时,BSON Decimal128 类型被转换为具有特定比例和精度的BigDecimal。能够指定规模和精度很重要,因为这将对集群的 RAM 和磁盘空间大小产生影响。
    • 查看以下第 309-313 行,了解 Mongo Spark 连接器在使用模式推断时如何推断 BigDecimal 比例和精度。 github.com/mongodb/mongo-spark/blob/master/src/main/scala/com/…
    • 作为对上述内容的澄清,我应该说“将数据从 Mongo 加载到 Spark 时,Spark SQL Schema 指定 Spark 如何表示数据,而不是数据如何存储在 Mongo 中。”
    • 对不起,我不知道。
    【解决方案2】:

    我可以通过绕过 load 辅助方法并直接在 MongoSpark 实例上调用 toDF(schema) 来做到这一点:

     val schema = StructType(
                                 List(StructField("time", LongType, false),
                                      StructField("pubKeyId", LongType, false),
                                      StructField("value", DecimalType(30, 0), false),
                                      StructField("outIndex", LongType, false),
                                      StructField("outTxId", LongType, false)
                                 ))
        val outputs =    
          builder().sparkContext(sc).readConfig(rc).build().toDF(schema).as[Output]
    

    这会产生正确的架构,并且数据被正确读入 Spark,没有任何缺失值:

        outputs.printSchema()
    
     |-- time: long (nullable = false)
     |-- pubKeyId: long (nullable = false)
     |-- value: decimal(30,0) (nullable = false)
     |-- outIndex: long (nullable = false)
     |-- outTxId: long (nullable = false)
    

    【讨论】:

    • 这种方法除了可以指定比例和范围之外,还有一个好处是我们可以指定该字段不可为空。这与 load[Output] 形成对比,load[Output] 始终使用 Decimal(38,18) 类型的可为空字段。
    猜你喜欢
    • 1970-01-01
    • 2019-07-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-04-04
    • 2017-02-14
    • 2016-05-27
    • 1970-01-01
    相关资源
    最近更新 更多