PySpark多级经理查询UDF报错求助:Method col([class java.lang.Integer])不存在
在Spark中获取员工指定层级内的所有上级经理问题解决
需求说明
我需要在Spark中实现获取员工指定层级内的所有上级经理的功能。示例数据集如下:
EMPLOYEE_ID,FIRST_NAME,LAST_NAME,EMAIL,PHONE_NUMBER,HIRE_DATE,JOB_ID,SALARY,COMMISSION_PCT,MANAGER_ID,DEPARTMENT_ID 1,Donald,OConnell,DOCONNEL,650.507.9833,21/06/2007,SH_CLERK,2600,,2,500 2,Douglas,Grant,DGRANT,650.507.9844,13/01/2008,SH_CLERK,2600,,3,50 3,Jennifer,Whalen,JWHALEN,515.123.4444,17/09/2003,AD_ASST,4400,,4,10 4,Michael,Hartstein,MHARTSTE,515.123.5555,17/02/2004,MK_MAN,13000,,5,20 5,Pat,Fay,PFAY,603.123.6666,17/08/2005,MK_REP,6000,,6,20 6,Susan,Mavris,SMAVRIS,515.123.7777,07/06/2002,HR_REP,6500,,7,40 7,Hermann,Baer,HBAER,515.123.8888,07/06/2002,PR_REP,10000,,8,70 8,Shelley,Higgins,SHIGGINS,515.123.8080,07/06/2002,AC_MGR,12008,,9,110 9,William,Gietz,WGIETZ,515.123.8181,07/06/2002,AC_ACCOUNT,8300,,,110
当查询员工ID为1时,预期结果为['3', '4', '5', '6', '7', '8', '9'](注:实际完整层级包含直接经理2,可按需调整逻辑)。
我尝试的PySpark代码(运行报错)
import sys import os from operator import add import re import pyspark.sql.functions as F os.environ['SPARK_HOME'] = "path" sys.path.append("path") sys.path.append("path") def recur_man(emp_id,lvl,list1): with open("path\\employee_1.txt") as f: for lines in f: if lines.split(',')[0] == emp_id: list1.append(lines.split(',')[9]) lvl-=1 recur_man(lines.split(',')[9],lvl,list1) return list1 try: from pyspark import SparkContext from pyspark import SparkConf from pyspark.sql import SQLContext from pyspark.sql.functions import * from pyspark.sql.types import * config = SparkConf().setAll([('spark.num.executors','10'),('spark.ui.port','4050')]) sc = SparkContext(conf=config) sqlContext = SQLContext(sc) rdd = sc.textFile("path\\employee_1.txt") header = rdd.first() header_mod = [x.encode("utf-8") for x in header.split(',')] rdd = rdd.filter(lambda line: line!=header) rdd = rdd.map(lambda line: line.split(',')) df1 = rdd.toDF(header_mod) spark_recur_man = udf(lambda x: recur_man,ArrayType(StringType())) list1 = [] df1.select('EMPLOYEE_ID',spark_recur_man('EMPLOYEE_ID',3,list1).alias('heirarchy')).show(truncate=False) sc.stop() except ImportError as e: print ("Error importing Spark Modules", e) sys.exit(1)
错误信息
df1.select('EMPLOYEE_ID',spark_recur_man('EMPLOYEE_ID',3,list1).alias('heirarchy')).show(truncate=False) File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\functions.py", line 1957, in wrapper File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\functions.py", line 1918, in call File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\column.py", line 60, in _to_seq File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\column.py", line 48, in _to_java_column File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\column.py", line 41, in _create_column_from_name File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\py4j-0.10.4-src.zip\py4j\java_gateway.py", line 1133, in __call__ File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\sql\utils.py", line 63, in deco File "C:\opt\spark\spark-2.2.0-bin-hadoop2.7\python\lib\py4j-0.10.4-src.zip\py4j\protocol.py", line 323, in get_return_value py4j.protocol.Py4JError: An error occurred while calling z:org.apache.spark.sql.functions.col. Trace: py4j.Py4JException: Method col([class java.lang.Integer]) does not exist
问题分析与解决方案
咱们先拆解下你代码里的核心问题,再给出正确的实现方式:
1. 直接导致报错的UDF问题
你在定义和调用UDF时犯了两个关键错误:
- UDF定义时,lambda只接收了一个参数,但你的
recur_man需要3个参数,而且你只是返回了函数对象,没有实际调用它 - 调用UDF时传入的
3是整数,Spark会把它当成列名去查找,这就是报错Method col([class java.lang.Integer]) does not exist的根源——Spark找不到名为3的列
2. 递归函数的设计缺陷
你的recur_man直接读取本地文件,这在Spark分布式环境中完全不可行:每个Executor都会独立读取本地文件,不仅重复加载数据,还可能因为节点间文件不一致导致结果错误。而且数据已经加载到DataFrame里了,完全没必要重复读文件。
正确实现:用Spark递归CTE处理层级关系
Spark 2.2+支持递归CTE(Common Table Expression),这是处理层级关系最简洁高效的方式,完全符合Spark的分布式计算模型。下面是完整的实现代码:
import sys import os from pyspark import SparkContext, SparkConf from pyspark.sql import SQLContext os.environ['SPARK_HOME'] = "path" sys.path.append("path") sys.path.append("path") try: config = SparkConf().setAll([('spark.num.executors','10'),('spark.ui.port','4050')]) sc = SparkContext(conf=config) sqlContext = SQLContext(sc) # 加载数据并转换为DataFrame rdd = sc.textFile("path\\employee_1.txt") header = rdd.first() header_list = header.split(',') rdd = rdd.filter(lambda line: line != header) rdd = rdd.map(lambda line: line.split(',')) df = rdd.toDF(header_list) # 注册为临时表,方便使用SQL查询 df.createOrReplaceTempView("employees") # 目标员工ID,可按需修改 target_emp_id = "1" # 递归CTE查询:获取所有上级经理 recursive_sql = f""" WITH RECURSIVE emp_hierarchy AS ( -- 基础查询:获取目标员工的直接经理 SELECT EMPLOYEE_ID, MANAGER_ID FROM employees WHERE EMPLOYEE_ID = '{target_emp_id}' UNION ALL -- 递归查询:迭代获取上一级经理的经理 SELECT e.EMPLOYEE_ID, e.MANAGER_ID FROM employees e JOIN emp_hierarchy eh ON e.EMPLOYEE_ID = eh.MANAGER_ID WHERE e.MANAGER_ID IS NOT NULL AND e.MANAGER_ID != '' ) -- 收集所有经理ID,排除空值,同时可按需过滤直接经理 SELECT COLLECT_LIST(MANAGER_ID) AS all_managers FROM emp_hierarchy WHERE MANAGER_ID IS NOT NULL AND MANAGER_ID != '' -- 如果需要排除直接经理(只取层级>=2的),添加下面的条件 -- AND MANAGER_ID != (SELECT MANAGER_ID FROM employees WHERE EMPLOYEE_ID = '{target_emp_id}') """ # 执行查询并展示结果 result_df = sqlContext.sql(recursive_sql) result_df.show(truncate=False) sc.stop() except ImportError as e: print("Error importing Spark Modules", e) sys.exit(1)
代码说明
- 递归CTE逻辑:先获取目标员工的直接经理,然后通过自连接不断向上迭代,直到没有上级经理为止
- 结果调整:如果需要排除直接经理(得到你预期的
['3','4','5','6','7','8','9']),只需打开代码中注释的过滤条件即可 - 性能优势:递归CTE是Spark原生支持的优化查询,比自定义UDF更高效,能充分利用Spark的分布式计算能力
内容的提问来源于stack exchange,提问作者user7012176
相关产品推荐
相关产品推荐

