【发布时间】:2018-04-21 03:59:02
【问题描述】:
我正在一个项目中使用 Scala 和 Spark 处理存储在 HDFS 中的文件。这些文件每天早上都会登陆 HDFS。我的工作是每天从 HDFS 读取该文件,对其进行处理,然后将结果写入 HDFS。在我将文件转换为 Dataframe 后,此作业执行过滤器以仅获取包含高于上一个文件中处理的最高时间戳的时间戳的行。此过滤器仅在几天内出现未知行为。尽管新文件包含与该过滤器匹配的行,但有些日子按预期工作,而有些日子则过滤结果为空。当它在 TEST 环境中执行时,同一文件总是发生这种情况,但在我的本地工作中,使用具有相同 HDFS 连接的同一文件按预期工作。
我尝试以不同的方式进行过滤,但是对于某些特定文件,没有一个可以在该环境中工作,但所有这些都可以在我的 LOCAL 中正常工作: 1)火花sql
val diff = fp.spark.sql("select * from curr " +
s"where TO_DATE(CAST(UNIX_TIMESTAMP(substring(${updtDtCol},
${substrStart},${substrEnd}),'${dateFormat}') as TIMESTAMP))" +
s" > TO_DATE(CAST(UNIX_TIMESTAMP('${prevDate.substring(0,10)}'
,'${dateFormat}') as TIMESTAMP))")
2) 火花过滤器函数
val diff = df.filter(date_format(unix_timestamp(substring(col(updtDtCol),0,10),dateFormat).cast("timestamp"),dateFormat).gt(date_format(unix_timestamp(substring(col("PrevDate"),0,10),dateFormat).cast("timestamp"),dateFormat)))
3) 使用过滤器的结果添加额外的列,然后按此新列过滤
val test2 = df.withColumn("PrevDate", lit(prevDate.substring(0,10)))
.withColumn("DatePre", date_format(unix_timestamp(substring(col("PrevDate"),0,10),dateFormat).cast("timestamp"),dateFormat))
.withColumn("Result", date_format(unix_timestamp(substring(col(updtDtCol),0,10),dateFormat).cast("timestamp"),dateFormat).gt(date_format(unix_timestamp(substring(col("PrevDate"),0,10),dateFormat).cast("timestamp"),dateFormat)))
.withColumn("x", when(date_format(unix_timestamp(substring(col(updtDtCol),0,10),dateFormat).cast("timestamp"),dateFormat).gt(date_format(unix_timestamp(substring(col("PrevDate"),0,10),dateFormat).cast("timestamp"),dateFormat)), lit(1)).otherwise(lit(0)))
val diff = test2.filter("x == 1")
我认为问题不是由过滤器本身引起的,也不是由文件引起的,但我希望收到有关我应该检查什么或是否有人以前遇到过此问题的反馈。
请让我知道在此处发布哪些信息可能有用以获得一些反馈。
文件示例的一部分如下所示:
|TIMESTAMP |Result|x|
|2017-11-30-06.46.41.288395|true |1|
|2017-11-28-08.29.36.188395|false |0|
将 TIMESTAMP 值与上一个日期(例如:2017-11-29)进行比较,然后我创建一个名为“结果”的列,该比较的结果始终适用于环境和另一个名为“x”的列结果相同。
正如我之前提到的,如果我在两个日期之间使用比较器函数或在“结果”或“x”列中使用结果来过滤数据帧,有时结果是一个空数据帧,但在本地使用相同的 HDFS 和文件,结果包含数据。
【问题讨论】:
标签: apache-spark spark-dataframe