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定义和实例方法使用存在三个关键问题:
- UDF定义时机错误:
ataccama_cleanse_udf作为类属性初始化时,self未绑定到JobBase实例,lambda中的self指向类本身而非实例,导致调用UDF时传递了错误对象类型。 - 实例方法缺少self参数:
__cleanse_dnb_attribute作为实例方法未声明self参数,无法访问实例的ataccama_url属性,同时引发参数传递混乱。 - 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)
关键修改点说明
- 调整UDF定义时机:将
ataccama_cleanse_udf的定义移到execute方法中,确保在ataccama_url赋值完成后再初始化UDF。 - 修复实例方法参数:给
__cleanse_dnb_attribute添加self参数,使其能正确访问实例属性。 - 修正DataFrame赋值逻辑:拆分
withColumn和show调用,避免show()返回的None覆盖原DataFrame变量。 - 统一返回值:将分支中的空字符串返回统一到方法末尾,避免未明确返回值的情况。
额外注意事项
- PySpark UDF运行在Worker节点上,需确保Worker节点能访问
ataccama_url对应的服务,必要时配置VPC或安全组规则。 - 如果Glue环境中缺少
requests库,需通过Glue Job的依赖包配置上传包含requests的wheel包。
内容的提问来源于stack exchange,提问作者pbh
相关产品推荐
相关产品推荐

