Spark Scala DataFrame交易表回溯逻辑实现问题
交易表回溯逻辑修正方案
原代码问题分析
- 连接别名错误:左连接条件里的
unf.P_ID是未定义的别名,实际应为tab1.P_ID - 时间范围计算偏差:使用
add_months是按月偏移,不符合“过去30天”的需求,应改用date_sub按天计算 - 存在性判断逻辑缺陷:左连接后用
distinct无法准确处理多匹配场景,直接的case when会因重复匹配导致结果异常 - 未过滤当日数据:原代码处理全量数据,不符合“按日加载2024-04-23数据”的要求
修正后的实现方案
方案1:使用EXISTS子查询(推荐,逻辑更清晰)
val tmpDf: DataFrame = spark.sqlContext.sql(""" SELECT t1.P_ID, t1.Desc, t1.P_Date, CASE WHEN EXISTS ( SELECT 1 FROM tab1 t2 WHERE t2.P_ID = t1.P_ID AND to_date(t2.P_Date) >= date_sub(to_date(t1.P_Date), 30) AND to_date(t2.P_Date) < to_date(t1.P_Date) ) THEN 'N' ELSE 'Y' END AS UNIQUE_FLAG FROM tab1 t1 WHERE to_date(t1.P_Date) = '2024-04-23' """)
方案2:使用左连接+聚合判断
val tmpDf: DataFrame = spark.sqlContext.sql(""" SELECT t1.P_ID, t1.Desc, t1.P_Date, CASE WHEN COUNT(t2.P_ID) > 0 THEN 'N' ELSE 'Y' END AS UNIQUE_FLAG FROM tab1 t1 LEFT JOIN tab1 t2 ON t1.P_ID = t2.P_ID AND to_date(t2.P_Date) >= date_sub(to_date(t1.P_Date), 30) AND to_date(t2.P_Date) < to_date(t1.P_Date) WHERE to_date(t1.P_Date) = '2024-04-23' GROUP BY t1.P_ID, t1.Desc, t1.P_Date """)
代码说明
- 先过滤出当日(2024-04-23)的数据,减少计算量
- 用
date_sub(to_date(t1.P_Date), 30)精准计算过去30天的起始日期 - EXISTS子查询直接判断该P_ID是否在过去30天有历史记录,效率更高
- 左连接+聚合的方式通过计数判断是否存在匹配记录,适合需要额外统计历史次数的场景
内容的提问来源于stack exchange,提问作者user12175004
相关产品推荐
相关产品推荐

