Scalatest中对比相同Spark Dataset断言失败问题
问题分析与解决方案
核心问题1:数据类型隐式转换不匹配
你的driveAggMonthly类中numberOfDrives字段是BigInt类型,但nr_drives_by_month函数里,count(col("driveId"))返回的是Long类型。虽然Spark的Schema显示两者都是bigint,但在Scala层面,Long和BigInt是不同的类型,这会导致Dataset的实际数据类型不兼容,进而触发断言失败。
核心问题2:Dataset直接用ScalaTest断言不可靠
ScalaTest的===操作符无法正确对比分布式的Dataset对象,即便数据和Schema看起来一致,底层的Dataset实例元数据(比如分区、执行计划)差异也会导致断言失败,必须将数据收集到本地集合后再对比。
修复步骤
1. 修正聚合函数的类型转换
修改nr_drives_by_month函数,将count的结果显式转换为BigInt:
def nr_drives_by_month(input_ds: Dataset[Drive]): Dataset[driveAggMonthly] = { input_ds .groupBy(month(to_date(col("date"))).as("Month")) // 将count的Long结果转为BigInt .agg(count(col("driveId")).cast(BigIntType).as("numberOfDrives")) .as[driveAggMonthly] }
或者在Scala层面做类型转换:
.agg(count(col("driveId")).as[Long].map(BigInt(_)).as("numberOfDrives"))
2. 正确的断言方式
不要直接对比Dataset对象,而是将两者的数据收集到本地Seq后再断言:
// 收集数据到本地集合 val expectedData = expected_ds_t1.collect().toSeq val resultData = result_ds_t1.collect().toSeq // 断言数据一致 assert(resultData === expectedData)
如果需要额外验证Schema一致性,可以添加:
assert(result_ds_t1.schema === expected_ds_t1.schema)
验证说明
- 类型修正后,Dataset中的
numberOfDrives字段会真正匹配BigInt类型,避免Scala层面的类型不兼容问题。 - 收集到本地后对比是分布式数据测试的标准做法,能绕过Dataset实例本身的元数据差异,只对比实际业务数据。
内容的提问来源于stack exchange,提问作者Vanessa Hydoski
相关产品推荐
相关产品推荐

