Replace datetime column in DF with hour from int column

Viewed 37

I have a pyspark df with an hour column (int) like this:

hour
0
0
1
...
14

And I have an execution_datetime variable that looks like 2022-01-02 17:23:11

Now I want to calculate a new column for my DF that holds my execution_datetime but the hour is replaced by the values from the hour col. Output should look like:

hour exec_dttm_with_hour
0 2022-01-02 00:23:11
0 2022-01-02 00:23:11
1 2022-01-02 01:23:11
... ...
14 2022-01-02 14:23:11

I know there are ways using i.e. .collect(), then edit the list and insert as new col. But I need to make use of sparks parallel execution since it could be a super high data load. Also, casting it to pandas and then editing it is not suitable for my use case.

Thanks in advance for any suggestions!

2 Answers

You can use the make_timestamp function if you're on spark 3+, else you can use an UDF.

spark.range(15). \
    withColumnRenamed('id', 'hour'). \
    withColumn('static_dttm', func.lit(execution_datetime).cast('timestamp')). \
    withColumn('dttm',
               func.expr('''make_timestamp(year(static_dttm), 
                                           month(static_dttm), 
                                           day(static_dttm), 
                                           hour, 
                                           minute(static_dttm), 
                                           second(static_dttm)
                                           )
                         ''')
               ). \
    drop('static_dttm'). \
    show()

# +----+-------------------+
# |hour|               dttm|
# +----+-------------------+
# |   0|2022-01-02 00:23:11|
# |   1|2022-01-02 01:23:11|
# |   2|2022-01-02 02:23:11|
# |   3|2022-01-02 03:23:11|
# |   4|2022-01-02 04:23:11|
# |   5|2022-01-02 05:23:11|
# |   6|2022-01-02 06:23:11|
# |   7|2022-01-02 07:23:11|
# |   8|2022-01-02 08:23:11|
# |   9|2022-01-02 09:23:11|
# |  10|2022-01-02 10:23:11|
# |  11|2022-01-02 11:23:11|
# |  12|2022-01-02 12:23:11|
# |  13|2022-01-02 13:23:11|
# |  14|2022-01-02 14:23:11|
# +----+-------------------+

Using UDF

def update_ts(string_ts, hour_col):
    import datetime

    dttm = datetime.datetime.strptime(string_ts, '%Y-%m-%d %H:%M:%S')

    return datetime.datetime(dttm.year, dttm.month, dttm.day, hour_col, dttm.minute, dttm.second)

update_ts_udf = func.udf(update_ts, TimestampType())

spark.range(15). \
    withColumnRenamed('id', 'hour'). \
    withColumn('dttm', update_ts_udf(func.lit(execution_datetime), func.col('hour'))). \
    show()

# +----+-------------------+
# |hour|               dttm|
# +----+-------------------+
# |   0|2022-01-02 00:23:11|
# |   1|2022-01-02 01:23:11|
# |   2|2022-01-02 02:23:11|
# |   3|2022-01-02 03:23:11|
# |   4|2022-01-02 04:23:11|
# |   5|2022-01-02 05:23:11|
# |   6|2022-01-02 06:23:11|
# |   7|2022-01-02 07:23:11|
# |   8|2022-01-02 08:23:11|
# |   9|2022-01-02 09:23:11|
# |  10|2022-01-02 10:23:11|
# |  11|2022-01-02 11:23:11|
# |  12|2022-01-02 12:23:11|
# |  13|2022-01-02 13:23:11|
# |  14|2022-01-02 14:23:11|
# +----+-------------------+

You can use concat to connect multiple strings in a row and lit to add a constant value to each row.

In the following code, a new column timestamp is introduced where the first 11 characters of execution_datetime are conceited with characters after the hours and in between the hours are added. It also makes sure, that the hours have a leading zero.

import pyspark.sql.functions as f

df = df.withColumn('timestamp', f.concat(f.lit(execution_datetime [0:11]), f.lpad(f.col('hour'), 2, '0') , f.lit(f.lit(execution_datetime [13:]))))

Remark: This might be the faster version than using timestamp functions as suggested in samkart's answer, but also less safe in catching wrong inputs.

Related