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

如何在Databricks PySpark中优化生存者规则应用(无需UDF)

多源数据合并优化需求(PySpark/Databricks)

我正在处理一个多源数据采集项目,现有一张包含80多万条记录的最终表。由于数据来源多样,不同数据源对各字段有优先级差异,因此维护了一张字段-数据源权重表(数值越小优先级越高),用于确定同字段不同数据源值的优先级;同时还有一张数据源-ID映射表(因其他用途无法省略)。

需要实现两种场景的合并逻辑:

  • 现有最终表的全量合并
  • 增量数据与现有记录的匹配合并

核心逻辑:按mobile/email分组合并不同数据源的记录,同组内根据权重表的优先级,选取各字段最高优先级的数据源值,最终每组生成一条合并记录。


参考表结构示例

1. 字段-数据源权重表

|   columns   | source1 | source2 | source3 | ... | source10 |
|-------------|---------|---------|---------|-----|----------|
| first_name  | 1       | 2       | 3       | ... | 10       |
| last_name   | 3       | 1       | 2       | ... | 9        |
| mobile      | 2       | 3       | 5       | ... | 8        |
| email       | 2       | 3       | 5       | ... | 8        |
| city        | 3       | 5       | 1       | ... | 9        |

2. 数据源-ID映射表

| source system | ID |
|---------------|----|
| source1       | 11 |
| source3       | 12 |
| source4       | 13 |
| source5       | 14 |
| source6       | 15 |

3. 最终表示例(清洗后约90个字段)

| first_name | last_name | mobile    | email          | city      | source system ID |
|------------|-----------|-----------|----------------|-----------|------------------|
| john       | doe       | 12345667  | john@xyz.com   | london    | 12               |
| j          |           | 12345667  | jack@abc.com   | new york  | 13               |
| john       | d         | 12345667  | john@abc.com   |           | 11               |
| jack       | dawn      | 344567754 | jack@abc.com   | amsterdam | 11               |
| j          | dawn      | 344567754 | jack@abc.com   | zurich    | 13               |

当前实现及性能问题

自定义UDF实现逻辑

目前通过自定义UDF实现字段优先级选择,示例代码如下:

def match_merge_surviver(col, col_input, inp_source_system_id, mobile):
    inp_weight = '999'
    temp = []
    first_name = {
        '17':1,
        '16':9,
        '14':8,
        '4':4,
        '13':3,
        '18':5,
        '19':2,
        '20':1,
        '21':1,
        '23':6,
        '11':999,
        '24':1
    }
    
    if(not col_input or col_input =='Not Provided' or col_input == 'Null'):
        return col_input
        
    if col =='first_name':
        temp={}
        try:
            if temp:
                if temp[mobile] > inp_weight:
                    temp[mobile] = inp_weight
                    return col_input
                else:
                    pass
            else:
                inp_weight = first_name[str(inp_source_system_id)]
                temp[mobile] = inp_weight
                return col_input
        except KeyError:
                # Key is not present
                inp_weight = 999
                pass

# Registering function
sqlContext.udf.register("survivorship", survivorship)

Spark SQL调用示例(仅展示first_name字段)

select a.email, a.source_system_id, a.mobile, a.DOB, a.first_name,
match_merge_surviver('first_name',a.first_name, cast(a.source_system_id as string),a.mobile) 
as sur_first_name 
from final_table a
where a.source_system_id <> 0 
and a.source_system_id is not null 
and a.mobile = '12345667'

性能瓶颈

该方案性能较差,后续还需结合ROW_NUMBER进行复杂关联:

ROW_NUMBER() OVER (PARTITION BY mobile ORDER BY survived_first_name, modified_date desc )

期望输出

每组生成一条合并后的记录,示例如下:

| first_name | last_name | mobile   | email        | city   | source system ID |
|------------|-----------|----------|--------------|--------|------------------|
| john       | doe       | 12345667 | john@xyz.com | london | 12               |

优化诉求

寻求避免使用UDF的优化方案,以提升Databricks PySpark中的数据处理性能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 05:50:49