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

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写入的额外优化

  1. 调整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")
    
  2. 减少Shuffle:检查作业里有没有没必要的groupBy、join,提前过滤无效数据,减少要写入的数据量。
  3. 用压缩格式:写Parquet或ORC时开压缩,能大幅减少IO量:
    spark.conf.set("spark.sql.parquet.compression.codec", "snappy")
    
  4. 调整EMR集群配置:给Executor加内存和CPU核心,减少Executor数量,提升单个节点的处理能力;开启动态资源分配。

临时排查步骤

  1. 看EMR的YARN日志,定位是Map阶段还是Reduce阶段卡住。
  2. 打开Spark UI(默认4040端口),看任务的执行时间、Shuffle数据量、UDF的耗时占比。
  3. 拿小数据集测试优化后的UDF,确认性能提升后再跑全量数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 18:48:41