Spark DataFrame转API替代SQL报错:Py4JException问题求助
问题描述
现有Spark DataFrame结构如下:
-root |-- ME_KE: string (nullable = true) |-- CSPD_CAT: string (nullable = true) |-- EFF_DT: string (nullable = true) |-- TER_DT: string (nullable = true) |-- CREATE_DTM: string (nullable = true) |-- ELIG_IND: string (nullable = true)
原使用Spark SQL的查询代码:
df=spark.read.format('csv').load(SourceFilesPath+"\\cutdetl.csv",infraSchema=True,header=True) df.createOrReplaceTempView("cutdetl") spark.sql(f"""select me_ke, eff_dt, ter_dt, create_dtm from cutdetl where (elig_ind = 'Y') and ((to_date('{start_dt}','dd-mon-yyyy') between eff_dt and ter_dt) or (eff_dt between to_date('{start_dt}','dd-mon-yyyy') and to_date('{end_dt}','dd-mon-yyyy'))) """)
(注:原SQL代码存在语法疏漏,已补全变量引号与闭合括号)
尝试转为DataFrame API实现时编写的代码:
df1=df.select("me_ke","eff_dt","ter_dt","elig_ind") .where(col("elig_ind")=="Y" & (F.to_date('31-SEP-2022', dd-mon-yyyy') .between(col("mepe_eff_dt"),col("mepe_term_dt"))) | (F.to_date(col("eff_dt")) .between(F.to_date('31-DEC-2022'),F.to_date('31-DEC-2022'))))
运行报错:
py4j.Py4JException: Method and([class java.lang.String]) does not exist
错误原因分析
- 逻辑运算符优先级错误:Python中
==优先级高于&/|,导致col("elig_ind")=="Y" & ...被错误解析为字符串与布尔值的位运算,触发类型不匹配错误。 - 列名误用:代码中错误使用
mepe_eff_dt/mepe_term_dt,实际列名为eff_dt/ter_dt。 - 参数格式缺失:日期格式字符串
dd-mon-yyyy未加引号,不符合字符串参数要求。 - 逻辑与原SQL不匹配:硬编码固定日期替代原SQL的
start_dt/end_dt变量,且部分to_date调用缺少格式参数。 - 语法不完整:代码存在括号未闭合问题,逻辑结构混乱。
修正后的DataFrame API代码
from pyspark.sql import functions as F from pyspark.sql.functions import col # 读取数据(与原代码保持一致) df = spark.read.format('csv').load(SourceFilesPath + "\\cutdetl.csv", infraSchema=True, header=True) # 转换为DataFrame API实现 df1 = df.select("me_ke", "eff_dt", "ter_dt", "create_dtm") \ .where( (col("elig_ind") == "Y") & ( (F.to_date(F.lit(start_dt), "dd-mon-yyyy").between(col("eff_dt"), col("ter_dt"))) | (col("eff_dt").between(F.to_date(F.lit(start_dt), "dd-mon-yyyy"), F.to_date(F.lit(end_dt), "dd-mon-yyyy"))) ) )
关键修正说明
- 括号明确逻辑顺序:将每个条件用括号包裹,确保
==判断与&/|逻辑运算顺序正确。 - 变量正确转换:用
F.lit()将Python变量start_dt/end_dt转为Spark字面量,匹配原SQL的变量逻辑。 - 修正列名与参数:替换错误列名,给所有
to_date调用补全格式字符串"dd-mon-yyyy"。 - 匹配原SQL字段选择:在
select中保留原SQL的create_dtm字段,确保输出一致。
内容的提问来源于stack exchange,提问作者venkat
相关产品推荐
相关产品推荐

