Delta分区表跨月/年场景下7天数据通用过滤方案咨询
解决Delta表跨年月分区的日期过滤问题
原代码的核心问题是将年、月、日三个分区字段独立做范围判断,当日期跨月/年时,日的范围条件会完全失效。比如触发日2024-05-07对应的起始日2024-04-29、结束日2024-05-05,day between 29 and 5会过滤掉所有记录,导致返回空数据。
以下是两种通用的解决方案,均可处理跨月、跨年的场景:
方案一:拼接分区字段为日期类型后比较
from datetime import datetime, timedelta from pyspark.sql.functions import col, to_date, concat_ws trigger_date = "2024-05-07" # yyyy-mm-dd格式 trigger_date_obj = datetime.strptime(trigger_date, '%Y-%m-%d') # 计算目标时间范围:上周周一至周日 start_date = (trigger_date_obj - timedelta(days=8)).strftime('%Y-%m-%d') end_date = (trigger_date_obj - timedelta(days=2)).strftime('%Y-%m-%d') # 将分区的year、month、day拼接为标准日期字符串,再转为日期类型 condition = to_date(concat_ws("-", col("year"), col("month"), col("day")), 'yyyy-MM-dd').between(start_date, end_date) data = spark.read.format("delta").table("delta_table_name").filter(condition)
该方案利用Spark的日期类型处理逻辑,自动兼容跨月/跨年的日期范围判断,无需手动处理年月的边界问题。
方案二:拼接为标准格式字符串做范围比较(性能更优)
from datetime import datetime, timedelta from pyspark.sql.functions import col, concat_ws, lpad trigger_date = "2024-05-07" # yyyy-mm-dd格式 trigger_date_obj = datetime.strptime(trigger_date, '%Y-%m-%d') start_date = (trigger_date_obj - timedelta(days=8)).strftime('%Y-%m-%d') end_date = (trigger_date_obj - timedelta(days=2)).strftime('%Y-%m-%d') # 补全month和day为两位数字,拼接成yyyy-MM-dd格式的字符串 partition_date_str = concat_ws("-", col("year"), lpad(col("month"), 2, '0'), lpad(col("day"), 2, '0')) # 利用字符串字典序做范围判断 condition = partition_date_str.between(start_date, end_date) data = spark.read.format("delta").table("delta_table_name").filter(condition)
由于yyyy-MM-dd格式的字符串按字典序排列与日期顺序完全一致,直接比较字符串范围可以避免日期类型转换的开销,性能更优。注意必须用lpad将单数字的月/日补零为两位,否则字符串比较会出错。
内容的提问来源于stack exchange,提问作者Matthew
相关产品推荐
相关产品推荐

