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

AWS Glue调用PySpark UDF报错:TypeError参数类型异常

AWS Glue Job自定义PySpark UDF参数类型不匹配问题解决

问题描述

在AWS Glue Job中调用自定义PySpark UDF时触发TypeError,完整代码及报错信息如下:

完整代码

import sys,os
import concurrent.futures
from concurrent.futures import *
import boto3
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from pyspark.context import SparkConf
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.dynamicframe import DynamicFrame
from datetime import datetime
from pyspark.sql.functions import array
from pyspark.sql.functions import sha2, concat_ws
from pyspark.sql.functions import  udf
from pyspark.sql.functions import StringType
import requests
import json


###############################

class JobBase(object):
    
    fair_scheduler_config_file= "fairscheduler.xml"
    rowAsDict={}
    listVendorDF=[]
    Oracle_Username=None
    Oracle_Password=None
    Oracle_jdbc_url=None
    futures=[]
    ataccama_url=None
    #all spark configuations can be passed in object in s3 bucket 
    ataccama_cleanse_udf=udf(lambda x:self.__cleanse_dnb_attribute(x),StringType() )
    
    
    
    def __cleanse_dnb_attribute(v_dnb_attr):
        payload = '{"in":{"src_org_name":"' + v_dnb_attr +'","sco_in":0,"exp_in":""}}'
        r = requests.post(self.ataccama_url, data=payload)  # response
        r.raise_for_status()
        if r is not None:
            r2 = json.loads(r.text)
            if r2['out'] is not None:
                r3 = r2['out']['cio_org_name'].replace(' ', '')
                return r3
            else:
                ''
        else:
            ''

  
    def __start_spark_glue_context(self):
        conf = SparkConf().setAppName("python_thread").set('spark.scheduler.mode', 'FAIR').set("spark.scheduler.allocation.file", self.fair_scheduler_config_file)
        self.sc = SparkContext(conf=conf)
        self.glueContext = GlueContext(self.sc)
        self.spark = self.glueContext.spark_session
        
        
    def __spark_read_from_table(self,table_name):
        return self.glueContext.read.format("jdbc").option("url", self.Oracle_jdbc_url).option("dbtable", table_name).option("user", self.Oracle_Username).option("password", self.Oracle_Password).option("numPartitions",2)\
        .option("lowerBound", 1)\
        .option("upperBound",10000)\
        .option("partitionColumn", "ORG_CODE").load()
        
            
    def execute(self):
        self.__start_spark_glue_context()
        
        args = getResolvedOptions(sys.argv, ['JOB_NAME','ataccma-cleanse-url'])
        self.ataccama_url=args['ataccma_cleanse_url']
        self.logger = self.glueContext.get_logger()
        self.logger.info("Starting Glue Threading job ")
        client = boto3.client('glue', region_name='XXXXXXXXXX')
        response = client.get_connection(Name='XXXXXXXX')
        connection_properties = response['Connection']['ConnectionProperties']
        URL = connection_properties['JDBC_CONNECTION_URL']
        url_list = URL.split("/")
        host = "{}".format(url_list[-2][:-5])
        new_host=host.split('@',1)[1]
        port = url_list[-2][-4:]
        database = "{}".format(url_list[-1])
        self.Oracle_Username = "{}".format(connection_properties['USERNAME'])
        self.Oracle_Password = "{}".format(connection_properties['PASSWORD'])
        spark_pool_configuration=3
        print("Host:",host)
        print("New Host:",new_host)
        print("Port:",port)
        print("Database:",database)
        self.Oracle_jdbc_url="jdbc:oracle:thin:@//"+new_host+":"+port+"/"+database
        print("Oracle_jdbc_url:",self.Oracle_jdbc_url)

        source_df =self.spark.read.format("jdbc").option("url", self.Oracle_jdbc_url).option("dbtable", "(select ENTERPRISE_NUM,ENTERPRISE_NAME,DNB_BUS_NM_TXT,DNB_SITE_BUS_STR_TXT  from xxgmdmadm.mdm_firmographic_data_v2 where ORG_ACCT_ID in (11758718960,11758836692)) ").option("user", self.Oracle_Username).option("password", self.Oracle_Password).load()
        source_df.show(truncate=False)
        
        source_df=source_df.withColumn('DNB_BUS_NM_TXT_CLEANSED',self.ataccama_cleanse_udf( source_df['DNB_BUS_NM_TXT'])).show(truncate=False)


def main():
    job = JobBase()
    job.execute() 


if __name__ == '__main__':
    main()

报错信息

TypeError: Invalid argument, not a string or column: <main.JobBase object at 0x7f4a77382390> of type <class 'main.JobBase'>. For column literals, use 'lit', 'array', 'struct' or 'create_map' function.

问题分析

报错核心是UDF定义和实例方法使用存在三个关键问题:

  1. UDF定义时机错误:ataccama_cleanse_udf作为类属性初始化时,self未绑定到JobBase实例,lambda中的self指向类本身而非实例,导致调用UDF时传递了错误对象类型。
  2. 实例方法缺少self参数:__cleanse_dnb_attribute作为实例方法未声明self参数,无法访问实例的ataccama_url属性,同时引发参数传递混乱。
  3. UDF绑定实例方法方式错误:类属性中定义UDF无法关联实例方法,因为实例属性(如ataccama_url)在类初始化阶段还未赋值。

解决方案

针对上述问题,修改代码如下:

修改后的核心代码

class JobBase(object):
    
    fair_scheduler_config_file= "fairscheduler.xml"
    rowAsDict={}
    listVendorDF=[]
    Oracle_Username=None
    Oracle_Password=None
    Oracle_jdbc_url=None
    futures=[]
    ataccama_url=None
    # 移除类属性中的UDF定义,改为在实例方法中初始化
    ataccama_cleanse_udf = None
    
    def __cleanse_dnb_attribute(self, v_dnb_attr):
        # 添加self参数,正确访问实例的ataccama_url
        payload = '{"in":{"src_org_name":"' + v_dnb_attr +'","sco_in":0,"exp_in":""}}'
        r = requests.post(self.ataccama_url, data=payload)
        r.raise_for_status()
        if r is not None:
            r2 = json.loads(r.text)
            if r2.get('out') is not None:
                r3 = r2['out']['cio_org_name'].replace(' ', '')
                return r3
        return ''  # 统一返回空字符串,避免分支返回None
    
    def execute(self):
        self.__start_spark_glue_context()
        
        args = getResolvedOptions(sys.argv, ['JOB_NAME','ataccma-cleanse-url'])
        self.ataccama_url=args['ataccma_cleanse_url']
        # 在实例初始化后定义UDF,此时self已绑定到实例
        self.ataccama_cleanse_udf = udf(lambda x: self.__cleanse_dnb_attribute(x), StringType())
        
        self.logger = self.glueContext.get_logger()
        self.logger.info("Starting Glue Threading job ")
        # 省略其余原有代码...
        
        source_df =self.spark.read.format("jdbc").option("url", self.Oracle_jdbc_url).option("dbtable", "(select ENTERPRISE_NUM,ENTERPRISE_NAME,DNB_BUS_NM_TXT,DNB_SITE_BUS_STR_TXT  from xxgmdmadm.mdm_firmographic_data_v2 where ORG_ACCT_ID in (11758718960,11758836692)) ").option("user", self.Oracle_Username).option("password", self.Oracle_Password).load()
        source_df.show(truncate=False)
        
        # 拆分withColumn和show,避免show返回的None覆盖原DataFrame
        source_df = source_df.withColumn('DNB_BUS_NM_TXT_CLEANSED', self.ataccama_cleanse_udf(source_df['DNB_BUS_NM_TXT']))
        source_df.show(truncate=False)

关键修改点说明

  1. 调整UDF定义时机:将ataccama_cleanse_udf的定义移到execute方法中,确保在ataccama_url赋值完成后再初始化UDF。
  2. 修复实例方法参数:给__cleanse_dnb_attribute添加self参数,使其能正确访问实例属性。
  3. 修正DataFrame赋值逻辑:拆分withColumn和show调用,避免show()返回的None覆盖原DataFrame变量。
  4. 统一返回值:将分支中的空字符串返回统一到方法末尾,避免未明确返回值的情况。

额外注意事项

  • PySpark UDF运行在Worker节点上,需确保Worker节点能访问ataccama_url对应的服务,必要时配置VPC或安全组规则。
  • 如果Glue环境中缺少requests库,需通过Glue Job的依赖包配置上传包含requests的wheel包。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 17:25:07