EMR上PySpark作业写入S3卡住,新增UDF后性能下降求助
PySpark EMR写入S3卡顿问题排查与优化
问题背景
在EMR集群跑PySpark作业,往S3写数据时,卡了30分钟没动静,只看到S3里的临时文件夹。已经加了这个参数:
spark_context.hadoopConfiguration.set("mapreduce.fileoutputcommitter.algorithm.version", "2")
但完全没用。之前作业跑得好好的,加了下面三个自定义UDF之后,性能直接崩了:
@udf(returnType=StringType()) def clean_phone(phone): try: if phone is not None: phone_cleaned= ''.join(e for e in phone if (e.isnumeric() or e=='+' or 'EXT' in phone.upper() )) if len(phone_cleaned)==11 and phone_cleaned.replace('+','').isdigit(): phone_number=str(phonenumbers.parse(phone_cleaned.replace("+",''),"US")).split(" ") if phone_number[5]==10 : phone_number[2]+"-"+phone_number[5] # 这里漏了return! else: return "" elif len(phone_cleaned)<10 : return "" else: phone_number=str(phonenumbers.parse(phone_cleaned,"US")).split(" ") if len(phone_number[5])==10 : return phone_number[2]+"-"+phone_number[5]+('-'+phone_number[7] if 'EXT' in phone_cleaned.upper() else "") else: return "" else: return "" except Exception as x: print("Error Occured in phone udf, Error: " + str(x)) @udf(returnType=StringType()) def clean_email(email): try: regex = r'\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Z|a-z]{2,}\b' if email is not None: if '%20' in email: email=email.replace("%20","") if '.jpg' in email: email="" email_corrected= ''.join(e for e in email if (e.isalnum() or e in ['.', '@','-','_'])) if(re.fullmatch(regex, email_corrected)): return email_corrected else: return "" else: return "" except Exception as x: print("Error Occured in email udf, Error: " + str(x)) @udf(returnType=StringType()) def clean_zip(zip): try: regex = r'[A-Za-z]' if zip is not None and len(re.findall('\d\d\d',zip))>0: zip_corrected=''.join(e for e in zip if (e.isnumeric() or e=='-') ) if "-" in zip_corrected: """Check cases when we have "-" in zip code <5""" if (zip_corrected.split("-")[0]).isdigit(): zip_corrected=zip_corrected.split("-")[0] elif (zip.split("-")[1]).isdigit(): zip_corrected=zip_corrected.split("-")[1] elif zip_corrected.isdigit() and len(zip_corrected)<5 : """Check cases when ;enth of zip code <5""" zip_corrected= zip_corrected.zfill(5) else : return "" if len(zipcodes.matching(zip_corrected))>0: # 这个外部查询是性能杀手 return zip_corrected else: return "" except Exception as x: print("Error Occured in zip udf, Error: " + str(x))
问题根源:UDF的低效写法拖垮了作业
1. Python UDF的固有开销
Python UDF要在PySpark里来回做JVM-Python的序列化/反序列化,比Spark原生函数慢好几倍,数据量大的时候这个开销会被无限放大。能不用就不用,能用原生函数或者Pandas UDF替代就赶紧换。
2. 单个UDF里的坑
(1)clean_phone的问题
- 有个分支漏写了return:
if phone_number[5]==10 : phone_number[2]+"-"+phone_number[5],这个分支会返回None,触发后续异常处理,平白增加额外开销。 - 把
phonenumbers.parse()的结果转成字符串再split,完全是绕远路——直接用phonenumbers的API拿号码组件就行,没必要拆字符串。 - 重复调用
phone.upper(),提前算一次复用就行。
优化后的版本:
import phonenumbers from phonenumbers import NumberParseException from pyspark.sql.functions import udf from pyspark.sql.types import StringType @udf(returnType=StringType()) def clean_phone(phone): try: if phone is None: return "" phone_upper = phone.upper() has_ext = 'EXT' in phone_upper # 只保留数字和+,EXT标记单独存 phone_cleaned = ''.join(e for e in phone if e.isnumeric() or e == '+') if len(phone_cleaned) < 10: return "" # 直接用API解析号码,避免字符串拆分 try: num_obj = phonenumbers.parse(phone_cleaned, "US") except NumberParseException: return "" if not phonenumbers.is_valid_number(num_obj): return "" country_code = str(num_obj.country_code) national_num = phonenumbers.format_number(num_obj, phonenumbers.PhoneNumberFormat.NATIONAL).replace("-", "") if len(national_num) != 10: return "" result = f"{country_code}-{national_num}" if has_ext: result += "-EXT" return result except Exception as x: print(f"Error in phone udf: {str(x)}") return ""
(2)clean_email的问题
- 每次调用都重新编译正则表达式,提前编译好复用能省不少时间。
- 字符串过滤用
''.join()效率低,换成正则替换更高效。
优化后的版本:
import re from pyspark.sql.functions import udf from pyspark.sql.types import StringType # 提前编译正则,避免重复编译 EMAIL_REGEX = re.compile(r'\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}\b') INVALID_CHARS = re.compile(r'[^A-Za-z0-9._%+-@-]') @udf(returnType=StringType()) def clean_email(email): try: if email is None: return "" email_processed = email.replace("%20", "") if '.jpg' in email_processed: return "" # 用正则替换过滤无效字符,比循环join快 email_corrected = INVALID_CHARS.sub('', email_processed) if EMAIL_REGEX.fullmatch(email_corrected): return email_corrected.lower() # 统一小写,增强数据一致性 else: return "" except Exception as x: print(f"Error in email udf: {str(x)}") return ""
(3)clean_zip的问题
zipcodes.matching()是致命性能瓶颈:如果这个函数需要查询外部API或远程数据库,每处理一条数据就发一次请求,数据量大的时候直接卡爆。必须把有效邮编提前加载到本地,用广播变量分发到各个Executor。- 正则和字符串处理逻辑冗余,简化后能减少不必要的计算。
优化后的版本:
import re from pyspark.sql.functions import udf from pyspark.sql.types import StringType # 提前把所有有效美国邮编加载到集合,然后广播到集群 # 示例:valid_zips = set(zipcodes.list_all()) # broadcast_valid_zips = spark.sparkContext.broadcast(valid_zips) ZIP_REGEX = re.compile(r'^\d{5}$') @udf(returnType=StringType()) def clean_zip(zip_str): try: if zip_str is None: return "" # 只保留数字和- zip_corrected = ''.join(e for e in zip_str if e.isnumeric() or e == '-') # 处理短邮编补0 if zip_corrected.isdigit(): if len(zip_corrected) < 5: zip_corrected = zip_corrected.zfill(5) elif len(zip_corrected) > 5: return "" # 处理带-的情况,只取前5位 elif '-' in zip_corrected: zip_part = zip_corrected.split('-')[0] zip_corrected = zip_part.zfill(5) if zip_part.isdigit() else "" # 从广播变量查有效性,避免远程调用 if ZIP_REGEX.match(zip_corrected) and zip_corrected in broadcast_valid_zips.value: return zip_corrected return "" except Exception as x: print(f"Error in zip udf: {str(x)}") return ""
3. 更彻底的优化:替换Python UDF
- 用Pandas UDF:把普通UDF改成
pandas_udf,用矢量化操作批量处理数据,性能比Python UDF高一个量级。 - 用Spark原生函数:能不用UDF就不用,比如邮箱、邮编的清洗,用Spark内置的
regexp_replace、substring、length等函数就能实现,完全避免JVM-Python交互开销。
S3写入的额外优化
- 调整Spark写入参数:
# 增大分区文件大小,减少小文件数量 spark.conf.set("spark.sql.files.maxPartitionBytes", "128m") # 用更高效的S3 committer spark.conf.set("spark.hadoop.fs.s3a.committer.name", "directory") spark.conf.set("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace") - 减少Shuffle:检查作业里有没有没必要的
groupBy、join,提前过滤无效数据,减少要写入的数据量。 - 用压缩格式:写Parquet或ORC时开压缩,能大幅减少IO量:
spark.conf.set("spark.sql.parquet.compression.codec", "snappy") - 调整EMR集群配置:给Executor加内存和CPU核心,减少Executor数量,提升单个节点的处理能力;开启动态资源分配。
临时排查步骤
- 看EMR的YARN日志,定位是Map阶段还是Reduce阶段卡住。
- 打开Spark UI(默认4040端口),看任务的执行时间、Shuffle数据量、UDF的耗时占比。
- 拿小数据集测试优化后的UDF,确认性能提升后再跑全量数据。
内容的提问来源于stack exchange,提问作者Xi12
相关产品推荐
相关产品推荐

