如何在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
相关产品推荐
相关产品推荐

