【问题标题】:Spark Structured streaming replacing values of a columnSpark结构化流替换列的值
【发布时间】: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-ddYYYYMMd

拼花写查询:

val tunerQuery = tunerDataJsonDF
  .writeStream
  .format("parquet")
  .option("path",pathtodata )
  .option("checkpointLocation", pathtochkpt)
  .partitionBy("date","env","serviceGroupId")
  .start()

【问题讨论】:

    标签: apache-spark spark-structured-streaming


    【解决方案1】:

    我假设您使用的是 Spark 2.2+

    tDataJsonDF.withColumn("formatted_date",date_format(to_date(col("date"), "YYYY/MM/dd"), "yyyy-MM-dd"))
    

    【讨论】:

      猜你喜欢
      • 2017-05-04
      • 2018-03-28
      • 2017-03-06
      • 1970-01-01
      • 2020-01-31
      • 1970-01-01
      • 1970-01-01
      • 2018-07-20
      • 1970-01-01
      相关资源
      最近更新 更多