mt8*_*t88 25 scala apache-spark
我在Spark中有一个数据框,有很多列和我定义的udf.我希望返回相同的数据帧,除非转换了一列.此外,我的udf接受一个字符串并返回一个时间戳.是否有捷径可寻?我试过了
val test = myDF.select("my_column").rdd.map(r => getTimestamp(r))
Run Code Online (Sandbox Code Playgroud)
但这会返回一个RDD,只返回已转换的列.
Dan*_*ula 41
如果你真的需要使用你的功能,我可以建议两个选项:
1)使用map/toDF:
import org.apache.spark.sql.Row
import sqlContext.implicits._
def getTimestamp: (String => java.sql.Timestamp) = // your function here
val test = myDF.select("my_column").rdd.map {
case Row(string_val: String) => (string_val, getTimestamp(string_val))
}.toDF("my_column", "new_column")
Run Code Online (Sandbox Code Playgroud)
2)使用UDF(UserDefinedFunction):
import org.apache.spark.sql.functions._
def getTimestamp: (String => java.sql.Timestamp) = // your function here
val newCol = udf(getTimestamp).apply(col("my_column")) // creates the new column
val test = myDF.withColumn("new_column", newCol) // adds the new column to original DF
Run Code Online (Sandbox Code Playgroud)
Bill Chambers在这篇很好的文章中有关于Spark SQL UDF的更多细节.
或者,
如果您只想将StringType列转换为TimestampType列,则可以使用自Spark SQL 1.5以来可用的unix_timestamp 列函数:
val test = myDF
.withColumn("new_column", unix_timestamp(col("my_column"), "yyyy-MM-dd HH:mm").cast("timestamp"))
Run Code Online (Sandbox Code Playgroud)
注意:对于火花1.5.x的,既要乘的结果unix_timestamp通过1000铸造时间戳之前(问题SPARK-11724).结果代码将是:
val test = myDF
.withColumn("new_column", (unix_timestamp(col("my_column"), "yyyy-MM-dd HH:mm") *1000L).cast("timestamp"))
Run Code Online (Sandbox Code Playgroud)
编辑:添加了udf选项
| 归档时间: |
|
| 查看次数: |
32105 次 |
| 最近记录: |