如何在DeltaLake Scala单元测试中创建带日期分区的Delta Lake布局?
如何在单元测试中创建带日期分区的Delta Lake结构
你可以直接构造带分区字段的测试数据,写入为分区的Delta Lake,替代原来转换已有Parquet的方式,具体步骤和代码如下:
步骤1:构造带分区列的测试数据
先创建包含ds分区字段的DataFrame,ds字段使用你需要的日期时间格式(比如2024-05-27 23:00:00),同时生成测试用的业务数据:
// 构造测试数据,包含分区列ds和业务字段 val testData = (1 to 42).map(i => (i, "2024-05-27 23:00:00")) .toDF("id", "ds")
步骤2:写入为分区的Delta Lake
将数据写入本地测试路径,指定按ds分区,Spark会自动生成ds=xxx的分区目录结构:
import org.apache.spark.sql.SaveMode // 本地测试路径(可复用你原有的路径变量) val deltaTestPath = singleInput.getParent.toFile.getAbsolutePath // 写入分区Delta Lake testData.write .mode(SaveMode.Overwrite) .format("delta") .partitionBy("ds") .save(deltaTestPath)
步骤3:读取并验证分区Delta Lake
用你已有的readDeltaLake方法读取数据,验证数据量和分区结构:
val myTestData = readDeltaLake(deltaTestPath) myTestData.count() shouldBe 42L // 可选:验证分区列是否正确识别 myTestData.select("ds").distinct().collect() should contain("2024-05-27 23:00:00")
关于分区目录的URL编码
你提到的ds=2024-05-27 23%3A00%3A00是URL编码后的格式,本地文件系统中Spark会直接生成ds=2024-05-27 23:00:00的目录,当通过HDFS等URI访问时会自动编码为%3A,不需要手动处理,读取时Spark能正常解析分区列。
替换原有代码后的完整测试片段
def readDeltaLake(path: String)(implicit sc: SparkSession): DataFrame = sc.read .format("org.apache.spark.sql.delta.sources.DeltaDataSource") .load(path) // local spark session implicit val sparkSession: SparkSession = aSparkSession() import sparkSession.implicits._ import org.apache.spark.sql.SaveMode // 构造带分区的测试数据 val testData = (1 to 42).map(i => (i, "2024-05-27 23:00:00")) .toDF("id", "ds") // 写入分区Delta Lake到测试路径 val deltaTestPath = singleInput.getParent.toFile.getAbsolutePath testData.write .mode(SaveMode.Overwrite) .format("delta") .partitionBy("ds") .save(deltaTestPath) // 读取并验证 val myTestData = readDeltaLake(deltaTestPath) myTestData.count() shouldBe 42L myTestData.select("ds").distinct().head().getString(0) shouldBe "2024-05-27 23:00:00"
内容的提问来源于stack exchange,提问作者Capacytron
相关产品推荐
相关产品推荐

