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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:15:45