【问题标题】:How to avoid boxing bytes in array in custom datasource?如何避免自定义数据源中数组中的装箱字节?
【发布时间】:2018-07-18 17:09:14
【问题描述】:

我正在开发自定义 Spark 数据源,并希望架构包含一行原始字节数组类型。

我的问题是生成的字节数组中的字节被装箱:然后输出类型为WrappedArray$ofRef。这意味着每个字节都表示为一个 java.lang.Object。虽然我可以解决这个问题,但我担心计算和内存开销,这对我的应用程序至关重要。我真的只想要原始数组!

以下是演示此行为的最小示例。

class DefaultSource extends SchemaRelationProvider with DataSourceRegister {

    override def shortName(): String = "..."

    override def createRelation(
                                    sqlContext: SQLContext,
                                    parameters: Map[String, String],
                                    schema: StructType = new StructType()
                               ): BaseRelation = {
        new DefaultRelation(sqlContext)
    }
}

class DefaultRelation(val sqlContext: SQLContext) extends BaseRelation with PrunedFilteredScan {

    override def schema = {
        StructType(
            Array(
                StructField("key", ArrayType(ByteType))
            )
        )
    }

    override def buildScan(
                              requiredColumnNames: Array[String],
                              filterArr: Array[Filter]
                          ): RDD[Row] = {
        testRDD
    }

    def testRDD = sqlContext.sparkContext.parallelize(
        List(
            Row(
                Array[Byte](1)
            )
        )
    )
}

使用此示例数据源如下:

def schema = StructType(Array(StructField("key", ArrayType(ByteType))))
val rows = sqlContext
        .read
        .schema(schema)
        .format("testdatasource")
        .load
        .collect()
println(rows(0)(0).getClass)

然后生成以下输出:

class scala.collection.mutable.WrappedArray$ofRef

在调试器中进一步检查结果类型确认 WrappedArray 中的字节确实被装箱 - 并且由于某种原因,它们的类型一直被擦除到java.lang.Object(而不是java.lang.Byte)。

请注意,直接使用 RDD,而不通过数据源 API,会导致原始字节数组的预期结果。

任何有关如何解决此问题的建议将不胜感激。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    好的,所以对于原始字节数组,我应该使用BinaryType 而不是Array(Byte) 作为列类型。这解决了问题。

    出于好奇,如果我们将 ArrayType(ByteType) 更改为例如ArrayType(LongType) 在上面的示例中,我们实际上得到了一个运行时异常,表明预期的装箱长。所以,Spark SQL 数组中的原语似乎总是被装箱的。

    【讨论】:

      猜你喜欢
      • 2013-07-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-09-05
      • 1970-01-01
      相关资源
      最近更新 更多