PySpark性能优化:移除耗时循环实现工程师分配逻辑
PySpark循环优化:批量分配优先级客户到区域可用工程师
问题背景
刚接触PySpark,现有一段耗时数小时的循环代码,希望移除循环优化性能。当前逻辑:遍历优先级客户DataFrame,检查工程师DataFrame中是否有客户所在区域的可用工程师;若存在可用工程师,则将该工程师的JOBS_LEFT减1,并标记客户已分配;若无则标记未分配。
原代码如下:
engAvail=spark.sql("SELECT EMPLOYEE_NO,concat('#',SECTOR,',') as SECTOR,JOBS_LEFT from tmpEmployeeAvailability WHERE JOBS_LEFT >0 ORDER BY JOBS_LEFT desc") for index,rows in df_priority_customer.iterrows(): engAvail = engAvail.filter(col("JOBS_LEFT")>0) pd = "#" + df_priority_customer.loc[index,'POSTAL_DISTRICT'] + "," empNo= engAvail.filter(col("SECTOR").contains(pd)) if (empNo.count() > 0) : #print(empNo.count()) #print(empNo.collect()[0][0]) #First Row engAvail = engAvail.withColumn('JOBS_LEFT', when(engAvail.EMPLOYEE_NO == int(empNo.collect()[0][0]), engAvail.JOBS_LEFT-1).otherwise(engAvail.JOBS_LEFT)) df_priority_customer.at[index,'ASSIGNED_ENGINEER'] = 'Y' else: #print('No Engineer') df_priority_customer.at[index,'ASSIGNED_ENGINEER'] = 'N'
原代码核心问题
- 用Pandas的
iterrows()遍历,将分布式的Spark作业退化为单节点循环,完全浪费Spark分布式计算优势 - 循环内频繁调用
count()、collect(),每次都会触发Spark作业提交,大量重复计算导致性能极低 - 逐行更新
engAvail和客户DataFrame,操作效率极低
优化方案
以下是基于Spark分布式特性的批量处理方案,全程避免循环:
1. 预处理数据,统一区域匹配格式
先将客户的邮编转换为与工程师区域一致的匹配格式,方便后续关联:
from pyspark.sql.functions import concat, lit, col df_priority_customer = df_priority_customer.withColumn( "MATCH_SECTOR", concat(lit("#"), col("POSTAL_DISTRICT"), lit(",")) )
2. 保留客户优先级顺序
原代码按客户DataFrame的行顺序遍历分配,因此需要给客户添加优先级序号,确保分配顺序与原逻辑一致:
from pyspark.sql import Window window_cust_order = Window.orderBy(lit(1)) # 按原DataFrame行顺序生成优先级 df_priority_customer = df_priority_customer.withColumn( "CUST_PRIORITY", row_number().over(window_cust_order) )
3. 关联客户与可用工程师
将客户与所在区域的可用工程师关联,并按客户优先级、工程师剩余任务量降序排序,保证优先分配剩余任务多的工程师:
cust_eng_join = df_priority_customer.join( engAvail, engAvail.SECTOR.contains(df_priority_customer.MATCH_SECTOR) & (engAvail.JOBS_LEFT > 0), "left" ).orderBy(df_priority_customer.CUST_PRIORITY, engAvail.JOBS_LEFT.desc())
4. 批量分配并标记客户状态
用窗口函数给每个客户匹配的工程师排序,取第一个可用工程师,同时标记客户是否分配成功:
window_assign = Window.partitionBy("CUST_PRIORITY").orderBy(engAvail.JOBS_LEFT.desc()) cust_assigned = cust_eng_join.withColumn( "ASSIGN_RANK", row_number().over(window_assign) ).filter(col("ASSIGN_RANK") == 1) # 生成最终客户分配结果 df_cust_result = cust_assigned.withColumn( "ASSIGNED_ENGINEER", when(col("EMPLOYEE_NO").isNotNull(), lit("Y")).otherwise(lit("N")) ).select(df_priority_customer.columns + ["ASSIGNED_ENGINEER", "EMPLOYEE_NO"])
5. 批量更新工程师剩余任务
计算每个工程师被分配的客户数量,一次性更新JOBS_LEFT:
from pyspark.sql.functions import sum # 统计每个工程师的分配次数 eng_assign_counts = cust_assigned.filter(col("EMPLOYEE_NO").isNotNull()) \ .groupBy("EMPLOYEE_NO") \ .agg(sum(lit(1)).alias("ASSIGN_COUNT")) # 更新工程师剩余任务量 updated_engAvail = engAvail.join( eng_assign_counts, on="EMPLOYEE_NO", how="left" ).withColumn( "JOBS_LEFT", when(col("ASSIGN_COUNT").isNotNull(), col("JOBS_LEFT") - col("ASSIGN_COUNT")).otherwise(col("JOBS_LEFT")) ).filter(col("JOBS_LEFT") > 0) # 过滤任务耗尽的工程师
优化效果
- 所有操作均为Spark分布式计算,避免单节点循环的性能瓶颈
- 仅触发有限次数的Spark作业,大幅减少重复计算开销
- 逻辑与原代码完全对齐:保持客户分配顺序、优先分配剩余任务多的工程师、正确更新工程师任务量
内容的提问来源于stack exchange,提问作者GaryH75
相关产品推荐
相关产品推荐

