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