基于Spark按指定班次规则拆分系统延迟数据的POC需求
问题内容整理
前置代码
from pyspark.sql.functions import current_timestamp, date_format, col, expr spark.createDataFrame( [ (1, "2024-10-19T13:01:17.000+00:00", "2024-10-21T13:01:17.000+00:00"), (2, "2024-10-10T11:01:17.000+00:00", "2024-10-10T12:01:17.000+00:00"), (3, "2024-10-11T06:00:00.000+00:00", "2024-10-12T05:59:59.000+00:00"), (4, "2024-10-21T14:00:00.000+00:00", "2024-10-21T15:16:07.000+00:00"), (5, "2024-10-24T04:00:00.000+00:00", "2024-10-24T11:00:00.000+00:00"), ], ["delay_id", "delay_start", "delay_end"], ).withColumn("delay_start_adjusted", expr("to_timestamp(delay_start)")).withColumn( "delay_end_adjusted", expr("to_timestamp(delay_end)") ).withColumn( "delay_start_adjusted_yyyymmdd", date_format(col("delay_start_adjusted"), "yyyyMMdd"), ).withColumn( "delay_end_adjusted_yyyymmdd", date_format(col("delay_end_adjusted"), "yyyyMMdd") ).withColumn( "delay_start_adjusted_hhmmss", date_format(col("delay_start_adjusted"), "HH:mm:ss") ).withColumn( "delay_end_adjusted_hhmmss", date_format(col("delay_end_adjusted"), "HH:mm:ss") ).createOrReplaceTempView( "temp_Delay" ) spark.createDataFrame( [ ("20241010",), ("20241011",), ("20241012",), ("20241019",), ("20241020",), ("20241021",), ("20241022",), ("20241023",), ("20241024",), ], ["date_key"], ).createOrReplaceTempView("temp_DateView")
需求说明
系统因各类原因产生延迟数据,每条延迟对应唯一delay_id。需完成以下POC任务:
- 基于上述生成的临时视图,按班次时间(0-8点、8-16点、16-24点)拆分延迟数据
- 最终输出的班次时间需减去6小时
- 必须使用
temp_DateView日期表进行关联
输出示例(以delay_id=3为例)
| date_key | delay_id | delay_start_adjusted_for_the_shift | delay_end_adjusted_for_the_shift | shift_code |
|---|---|---|---|---|
| 20241011 | 3 | 2024-10-11T00:00:00.000+00:00 | 2024-10-11T07:59:59.000+00:00 | S1 |
| 20241011 | 3 | 2024-10-11T08:00:00.000+00:00 | 2024-10-11T15:59:59.000+00:00 | S2 |
| 20241011 | 3 | 2024-10-11T16:00:00.000+00:00 | 2024-10-11T23:59:59.000+00:00 | S3 |
内容的提问来源于stack exchange,提问作者santoshkumar
相关产品推荐
相关产品推荐

