【发布时间】:2018-04-26 19:21:14
【问题描述】:
我有以下数据框
val tDataJsonDF = kafkaStreamingDFParquet
.filter($"value".contains("tUse"))
.filter($"value".isNotNull)
.selectExpr("cast (value as string) as tdatajson", "cast (topic as string) as env")
.select(from_json($"tdatajson", schema = ParquetSchema.tSchema).as("data"), $"env".as("env"))
.select("data.*", "env")
.select($"date", <--YYYY/MM/dd
$"time",
$"event",
$"serviceGroupId",
$"userId",
$"env")
此流数据框有一个日期列,其格式为 - YYYY/MM/dd。
因此,当我在 parquet write 中将此列用作分区列时,Spark 会将分区创建为date=2018%04%12。
有没有办法我可以在上面的代码中动态修改列值,使日期值为YYYY-MM-dd 或YYYYMMd。
拼花写查询:
val tunerQuery = tunerDataJsonDF
.writeStream
.format("parquet")
.option("path",pathtodata )
.option("checkpointLocation", pathtochkpt)
.partitionBy("date","env","serviceGroupId")
.start()
【问题讨论】:
标签: apache-spark spark-structured-streaming