如何用SQL与PySpark从秒数列生成日期和小时列(大数据场景)
Got it! Since you're dealing with large datasets where Pandas isn't feasible, PySpark's built-in date/time functions are perfect here—no need for slow UDFs. Let's walk through two approaches: using PySpark's DataFrame API directly, and using SQL queries, both of which will match your Pandas output exactly.
1. PySpark DataFrame API Implementation
We'll leverage Spark's optimized date functions to calculate the Date and hour columns without converting to Pandas:
from pyspark.sql import functions as F # Define our base start time (the reference date for the Time column) base_date = F.lit('2019-01-01') # Convert base date to Unix timestamp (seconds since 1970-01-01) base_unix_ts = F.unix_timestamp(base_date, 'yyyy-MM-dd') # Generate the Date and hour columns result_df = df.withColumn( 'Date', F.to_timestamp(base_unix_ts + F.col('Time')) # Add Time seconds to base timestamp, convert to readable date ).withColumn( 'hour', F.hour(F.col('Date')) # Extract hour from the Date timestamp ) # View the final result result_df.show(truncate=False)
Output:
+--------+-------------------+----+ |Time |Date |hour| +--------+-------------------+----+ |10.0 |2019-01-01 00:00:10|0 | |61.0 |2019-01-01 00:01:01|0 | |3500.0 |2019-01-01 00:58:20|0 | |3600.0 |2019-01-01 01:00:00|1 | |3700.54 |2019-01-01 01:01:40|1 | |7000.22 |2019-01-01 01:56:40|1 | |7200.22 |2019-01-01 02:00:00|2 | |15000.55|2019-01-01 04:10:00|4 | |86400.22|2019-01-02 00:00:00|0 | +--------+-------------------+----+
2. SQL Query Implementation
If you prefer working with SQL, you can register your DataFrame as a temporary view and run a SQL query with the same logic:
# Register the DataFrame as a temporary SQL view df.createOrReplaceTempView('time_events') # Execute the SQL query to generate the required columns sql_result = spark.sql(""" SELECT Time, to_timestamp(unix_timestamp('2019-01-01', 'yyyy-MM-dd') + Time) AS Date, hour(to_timestamp(unix_timestamp('2019-01-01', 'yyyy-MM-dd') + Time)) AS hour FROM time_events """) # View the SQL result sql_result.show(truncate=False)
This will produce the exact same output as the DataFrame API method.
Key Notes:
- Both approaches use Spark's native optimized functions, which are designed to handle large datasets efficiently (no memory issues like Pandas).
- If you need to change the base reference date, just replace
'2019-01-01'with your desired date in either method. - The
to_timestampfunction handles fractional seconds correctly, just like your Pandas implementation.
内容的提问来源于stack exchange,提问作者BhishanPoudel
相关产品推荐
相关产品推荐

