You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark SQL基于另一DataFrame日期值筛选交易数据

Spark按用户调研日期筛选交易记录实现方案

核心逻辑

调研表中person_id是用户唯一标识、取值全局唯一,不需要额外处理重复调研数据,只需要通过用户ID关联两张表,过滤出交易日期大于等于对应用户调研日期的交易记录即可,不需要复杂的窗口函数或者子查询。

注意:如果日期字段存储为字符串,且格式统一为yyyy-MM-dd,字符串字典序比较结果和日期比较结果一致;为了稳妥建议统一转为日期类型处理,避免不同日期格式导致的判断错误。

PySpark DataFrame API 实现

完整可运行代码如下:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date

# 初始化SparkSession
spark = SparkSession.builder.appName("filter_trans_by_survey").getOrCreate()

# 测试数据
survey_data = [["1","2022-03-11"], 
["2","2022-02-09"], 
["3","2022-01-21"], 
["4","2022-04-16"]]

transacction_Data =[[ "1", "2022-03-10"],
[ "1", "2022-03-14"],
[ "1", "2022-02-11"],
[ "2", "2022-01-30"],
[ "2", "2022-03-07"],
[ "2", "2022-02-16"],
[ "3", "2022-03-02"],
[ "4", "2022-05-15"]]

# 转换为DataFrame,同时将日期字段转为日期类型
survey_df = spark.createDataFrame(survey_data, schema=["person_id", "survey_date"]) \
    .withColumn("survey_date", to_date(col("survey_date"), "yyyy-MM-dd"))

trans_df = spark.createDataFrame(transacction_Data, schema=["person_id", "trans_date"]) \
    .withColumn("trans_date", to_date(col("trans_date"), "yyyy-MM-dd"))

# 关联+筛选逻辑
result_df = trans_df.join(
    survey_df,
    on="person_id",  # 按用户ID内连接
    how="inner"
).filter(
    col("trans_date") >= col("survey_date")  # 过滤交易日期晚于等于调研日期的记录
).select(
    col("person_id"), col("trans_date")  # 只保留需要输出的字段
)

# 查看结果
result_df.show()

运行后输出结果符合日期筛选规则:

+---------+----------+
|person_id|trans_date|
+---------+----------+
|        1|2022-03-14|
|        2|2022-02-16|
|        2|2022-03-07|
|        3|2022-03-02|
|        4|2022-05-15|
+---------+----------+

注:你提供的预期结果中遗漏了person_id=3的符合条件记录,该用户调研日期为2022-01-21,交易日期2022-03-02满足大于等于调研日期的规则,会被正常保留。

Spark SQL 实现

如果直接写SQL实现,逻辑相同,代码如下:

-- 关联筛选查询
SELECT 
    t.person_id,
    t.trans_date
FROM trans_table t
INNER JOIN survey_table s 
ON t.person_id = s.person_id
WHERE to_date(t.trans_date, 'yyyy-MM-dd') >= to_date(s.survey_date, 'yyyy-MM-dd');

生产环境注意事项

  • 如果后续业务调整为一个用户可能存在多条调研记录,需要先取每个用户最新的调研日期再做关联,避免关联后数据膨胀
  • 日期字段强烈建议转为日期/时间戳类型存储,不要用字符串存储,避免格式不一致导致的比较错误
  • 数据量较大的情况下,确保person_id字段作为join key,提前过滤无效数据,避免数据倾斜

内容的提问来源于stack exchange,提问作者Enzo Hernandez

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 10:03:22