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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 03:35:23