【发布时间】:2019-05-30 16:13:25
【问题描述】:
我正在创建一个 time_interval 列并将其添加到现有的 Pyspark 中的数据框。理想情况下,time_interval 将采用“HHmm”格式,分钟向下舍入到最接近的 15 分钟标记(815、830、845、900 等)。
我有为我执行逻辑的 spark sql 代码,但我如何获取连接为字符串列的值并将其插入现有数据帧?
time_interval = sqlContext.sql("select extract(hour from current_timestamp())||floor(extract(minute from current_timestamp())/15)*15")
time_interval.show()
+-------------------------------------------------------------------------------------------------------------------------------------------------------------------+
|concat(CAST(hour(current_timestamp()) AS STRING), CAST((FLOOR((CAST(minute(current_timestamp()) AS DOUBLE) / CAST(15 AS DOUBLE))) * CAST(15 AS BIGINT)) AS STRING))|
+-------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| 1045|
+-------------------------------------------------------------------------------------------------------------------------------------------------------------------+
baseDF = sqlContext.sql("select * from test_table")
newBase = baseDF.withColumn("time_interval", lit(str(time_interval)))
newBase.select("time_interval").show()
+--------------------+
| time_interval|
+--------------------+
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
|DataFrame[concat(...|
+--------------------+
only showing top 20 rows
因此,实际的预期结果应该只是在我正在创建的新列中显示实际的字符串值,而不是来自数据帧的连接值。如下所示:
newBase.select("time_interval").show(1)
+-------------+
|time_interval|
+-------------+
| 1045 |
+-------------+
【问题讨论】:
-
试试这个:
newBase = baseDF.selectExpr("*, extract(hour from current_timestamp())||floor(extract(minute from current_timestamp())/15)*15 AS time_interval") -
感谢 pault,“selectExpr” 的作用就像一个魅力!
标签: python apache-spark hive pyspark apache-spark-sql