Spark动态创建DataFrame时触发ArrayIndexOutOfBoundsException问题排查
问题描述
手动指定列和值创建DataFrame可正常运行,但通过遍历列动态生成数据时,执行testDf.show(false)会触发数组越界异常。
手动实现(正常运行)
代码
val columnSufix: String = "isNull" val data = Seq(Row( details.filter(col("DAY").isNull).count(), details.filter(col("CHANNEL_CATEGORY").isNull).count(), details.filter(col("SOURCE").isNull).count(), details.filter(col("PLATFORM").isNull).count() ) ) val schema: StructType = new StructType() .add(s"DAY_$columnSufix", LongType) .add(s"CHANNEL_CATEGORY_$columnSufix", LongType) .add(s"SOURCE_$columnSufix", LongType) .add(s"PLATFORM_$columnSufix", LongType) val testDf: DataFrame = spark.createDataFrame(spark.sparkContext.parallelize(data), schema) testDf.show(false)
执行输出
columnSufix: String = isNull data: Seq[org.apache.spark.sql.Row] = List([0,0,0,83845]) schema: org.apache.spark.sql.types.StructType = StructType(StructField(DAY_isNull,LongType,true),StructField(CHANNEL_CATEGORY_isNull,LongType,true),StructField(SOURCE_isNull,LongType,true),StructField(PLATFORM_isNull,LongType,true)) testDf: org.apache.spark.sql.DataFrame = [DAY_isNull: bigint, CHANNEL_CATEGORY_isNull: bigint ... 2 more fields]
动态实现(触发异常)
代码
val cols = details.columns.toSeq.take(4) val columnSuffix: String = "ISNULL" val data = cols.map(column => Row(details.filter(col(column).isNull).count())).toList val schema = StructType(cols.map(column => StructField(column + s"_$columnSuffix", LongType))) val testDf: DataFrame = spark.createDataFrame(spark.sparkContext.parallelize(data), schema) testDf.show(false)
报错信息
Caused by: RuntimeException: Error while encoding: java.lang.ArrayIndexOutOfBoundsException
if (assertnotnull(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else validateexternaltype(getexternalrowfield(assertnotnull(input[0, org.apache.spark.sql.Row, true]), 0, DAY_ISNULL), LongType, false) AS DAY_ISNULL#243076L
if (assertnotnull(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else validateexternaltype(getexternalrowfield(assertnotnull(input[0, org.apache.spark.sql.Row, true]), 1, CHANNEL_CATEGORY_ISNULL), LongType, false) AS CHANNEL_CATEGORY_ISNULL#243077L
if (assertnotnull(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else validateexternaltype(getexternalrowfield(assertnotnull(input[0, org.apache.spark.sql.Row, true]), 2, SOURCE_ISNULL), LongType, false) AS SOURCE_ISNULL#243078L
if (assertnotnull(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else validateexternaltype(getexternalrowfield(assertnotnull(input[0, org.apache.spark.sql.Row, true]), 3, PLATFORM_ISNULL), LongType, false) AS PLATFORM_ISNULL#243079L
Caused by: ArrayIndexOutOfBoundsException:
问题原因
动态代码中,data是由多个单元素Row组成的列表(例如List(Row(0), Row(0), Row(0), Row(83845))),但schema定义的是包含4个字段的StructType。Spark解析时期望每个Row包含4个字段匹配schema,实际每个Row仅1个元素,导致访问索引1、2、3时触发数组越界。
而手动实现中,data是包含一个4元素Row的Seq(List(Row(0,0,0,83845))),与schema的4个字段完全匹配,因此正常运行。
解决方法
方法1:合并统计值到单个Row
将所有动态生成的统计值包装进同一个Row,而非每个值单独生成Row:
val cols = details.columns.toSeq.take(4) val columnSuffix: String = "ISNULL" // 先收集所有统计值,再包装成单个Row val counts = cols.map(column => details.filter(col(column).isNull).count()) val data = Seq(Row.fromSeq(counts)) val schema = StructType(cols.map(column => StructField(column + s"_$columnSuffix", LongType))) val testDf: DataFrame = spark.createDataFrame(spark.sparkContext.parallelize(data), schema) testDf.show(false)
方法2:使用agg批量统计(性能更优)
避免多次filter操作,通过agg一次性统计所有列的null数量:
val cols = details.columns.toSeq.take(4) val columnSuffix: String = "ISNULL" // 批量生成聚合表达式,一次计算所有列的null数量 val aggExprs = cols.map(colName => count(when(col(colName).isNull, colName)).alias(s"${colName}_$columnSuffix")) val testDf = details.agg(aggExprs.head, aggExprs.tail:_*) testDf.show(false)
内容的提问来源于stack exchange,提问作者ultraInstinct

