Scala转PySpark:解析DataFrame where子句逻辑并完成代码转换
Scala Spark代码解析与PySpark转换
原代码where子句的作用
原Scala代码中的where子句用于筛选数据行:仅保留EVENT字段的值大于等于CONTINOUS_ENROL_START字段值,且同时小于等于CONTINOUS_ENROL_END字段值的记录——简单来说就是留下事件发生时间处于连续参保起止区间内的数据。
转换后的PySpark代码
# 需先导入col函数 from pyspark.sql.functions import col df = df1.join(df2, on=["ID"]).where((col("EVENT") >= col("CONTINOUS_ENROL_START")) & (col("EVENT") <= col("CONTINOUS_ENROL_END")))
转换说明
- PySpark中需要显式导入
col函数来引用DataFrame的字段; - 逻辑与操作符使用
&替代Scala中的&&,且每个条件需用括号包裹,避免运算优先级问题; join方法的关联字段参数用on=["ID"],和Scala的Seq("ID")作用一致。
内容的提问来源于stack exchange,提问作者user14269252
相关产品推荐
相关产品推荐

