【发布时间】:2018-07-01 11:34:15
【问题描述】:
如何在粘合作业中添加当前时间戳(额外列),以便输出数据具有额外列。在这种情况下:
架构源表: Col1, Col2
胶水作业后。
目的地架构: Col1, Col2, Update_Date(当前时间戳)
【问题讨论】:
标签: amazon-web-services pyspark etl aws-glue
如何在粘合作业中添加当前时间戳(额外列),以便输出数据具有额外列。在这种情况下:
架构源表: Col1, Col2
胶水作业后。
目的地架构: Col1, Col2, Update_Date(当前时间戳)
【问题讨论】:
标签: amazon-web-services pyspark etl aws-glue
我不确定DynamicFrame 是否有胶水原生方式来执行此操作,但您可以轻松转换为 Spark Dataframe,然后使用withColumn 方法。您需要使用lit 函数将文字值放入新列中,如下所示。
from datetime import datetime
from pyspark.sql.functions import lit
glue_df = glueContext.create_dynamic_frame.from_catalog(...)
spark_df = glue_df.toDF()
spark_df = spark_df.withColumn('some_date', lit(datetime.now()))
一些参考资料:
【讨论】:
from datetime import datetime 上面的代码 sn-p 中缺少语句
根据我使用 Glue 的经验,Glue 运行的时区是 GMT。但我的时区是 CDT。因此,要获得 CDT 时区,我需要在 SparkContext 中转换时间。这个具体的案例是添加last_load_date到target/sink。
所以我创建了一个函数。
def convert_timezone(sc):
sqlContext = SQLContext(sc)
local_time=dt.now().strftime('%Y-%m-%d %H:%M:%S')
local_time_df=sqlContext.createDataFrame([(local_time,)],['time'])
CDT_time_df = local_time_df.select(from_utc_timestamp(local_time_df['time'],'CST6CDT').alias('cdt_time'))
CDT_time=[i['cdt_time'].strftime('%Y-%m-%d %H:%M:%S') for i in CDT_time_df.collect()][0]
return CDT_time
然后像...一样调用函数
job_run_time = date_config.convert_timezone(sc)
datasourceDF0 = datasource0.toDF()
datasourceDF1 = datasourceDF0.withColumn('last_updated_date',lit(job_run_time))
【讨论】:
使用 Spark 的current_timestamp() 函数:
import org.apache.spark.sql.functions._
...
val timestampedDf = source.toDF().withColumn("Update_Date", current_timestamp())
val timestamped = DynamicFrame(timestampedDf, glueContext)
【讨论】:
正如我所见,这个问题没有正确的答案,我将尝试解释我对这个问题的解决方案:
首先要澄清 withColumn 函数是一个很好的方法,但重要的是要提到这个函数来自 Spark 本身的 Dataframe 并且这个函数不是胶水 DynamicFrame 的一部分,它是来自 Glue AWS 的自己的库,因此您需要隐藏框架来执行此操作....
第一步是从 DynamicFrame 获取 Spark Dataframe,胶水库使用函数 toDF() 函数完成此操作,一旦使用 Spark 框架,您就可以添加列和/或执行您需要的任何操作。
那么我们glue期望的是他自己的frame,所以我们需要从spark转回glue专有frame,为此可以使用DynamicFrame的apply函数,这需要导入对象:
import com.amazonaws.services.glue.DynamicFrame
并使用您应该已经拥有的glueContext,例如:
DynamicFrame(sparkDataFrame, glueContext)
在简历中,代码应如下所示:
import org.apache.spark.sql.functions._
import com.amazonaws.services.glue.DynamicFrame
...
val sparkDataFrame = datasourceToModify.toDF().withColumn("created_date", current_date())
val finalDataFrameForGlue = DynamicFrame(sparkDataFrame, glueContext)
...
注意:import org.apache.spark.sql.functions._ 是带上current_date() 功能添加日期列。
希望这会有所帮助....
【讨论】:
我们执行以下操作,并且无需转换为 DF() 就可以很好地工作
datasource0 = glueContext.create_dynamic_frame.from_catalog(...)
from datetime import datetime
def AddProcessedTime(r):
r["jobProcessedDateTime"] = datetime.today() #timestamp of when we ran this.
return r
mapped_dyF = Map.apply(frame = datasource0, f = AddProcessedTime)
【讨论】:
r 是什么?看起来像在 AddProcessedTime 函数中定义的参数,但是当您定义 f = AddProcessedTime 时,您没有将其作为参数传递?
您现在可以使用内置功能执行此操作:请参阅here...
请注意仅查找 glueContext.add_ingestion_time_columns 部分
【讨论】: