ORBIS企业多重复键值数据的通用整合策略咨询
ORBIS企业数据主键多行合并通用策略(附Pandas/Spark 3.5实现)
场景与需求
处理ORBIS企业数据时,存在同一主键(如company_id)对应多行数据的情况,包含三类数据:
- 完全重复的行(键值+字段内容均一致)
- 同一主键下语义相似但字段有差异的行(需保留多值)
- 同主键补充字段或异键的新数据
目标是将同一主键的所有有效数据合并为单行,既要去除冗余重复,又要保留同一字段的多值信息(避免仅保留首个值的局限)。
通用整合策略
- 预处理去重:先过滤完全重复的行,减少后续计算量
- 字段类型适配聚合:针对不同字段类型定义专属聚合规则:
- 字符串/枚举型:收集唯一值后用分隔符拼接,或保留所有有效语义值
- 数值型:根据业务需求选择保留所有唯一值、求和、均值或最值
- 固定属性字段(如行业、成立日期):取首个值或统一值
- 自定义逻辑扩展:针对语义相似的文本字段,可加入模糊匹配、同义词映射等预处理逻辑,再进行聚合
Pandas 可复现代码
构造测试数据
import pandas as pd # 模拟ORBIS企业数据 data = { 'company_id': ['C001', 'C001', 'C001', 'C002', 'C002'], 'company_name': ['ABC Corp', 'ABC Corp', 'ABC Limited', 'XYZ Inc', 'XYZ Inc'], 'industry': ['Manufacturing', 'Manufacturing', 'Manufacturing', 'Tech', 'Tech'], 'revenue': [1000000, 1000000, 1200000, 500000, 600000], 'establish_date': ['2000-01-01', '2000-01-01', '2000-01-01', '2010-05-10', '2010-05-10'] } df = pd.DataFrame(data)
实现合并逻辑
# 自定义聚合函数 def agg_unique_string(col): # 去重后用分号拼接唯一值 return '; '.join(col.unique()) def agg_unique_numeric(col): # 保留所有唯一数值,单个值则返回原值而非列表 unique_vals = col.unique().tolist() return unique_vals if len(unique_vals) > 1 else unique_vals[0] # 按主键分组,应用聚合规则 agg_rules = { 'company_name': agg_unique_string, 'industry': 'first', # 行业字段通常一致,取首个值 'revenue': agg_unique_numeric, 'establish_date': 'first' } merged_df = df.groupby('company_id').agg(agg_rules).reset_index() print(merged_df)
Spark 3.5 可复现代码
构造测试数据
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ val spark = SparkSession.builder() .appName("ORBISDataMerge") .master("local[*]") .getOrCreate() // 模拟ORBIS企业数据 val data = Seq( ("C001", "ABC Corp", "Manufacturing", 1000000L, "2000-01-01"), ("C001", "ABC Corp", "Manufacturing", 1000000L, "2000-01-01"), ("C001", "ABC Limited", "Manufacturing", 1200000L, "2000-01-01"), ("C002", "XYZ Inc", "Tech", 500000L, "2010-05-10"), ("C002", "XYZ Inc", "Tech", 600000L, "2010-05-10") ) val schema = StructType(Seq( StructField("company_id", StringType), StructField("company_name", StringType), StructField("industry", StringType), StructField("revenue", LongType), StructField("establish_date", StringType) )) val df = spark.createDataFrame(data).toDF(schema.fieldNames: _*)
实现合并逻辑(内置函数版,无需自定义UDF)
// 利用Spark 3.5内置函数实现去重+聚合 val mergedDF = df.groupBy("company_id") .agg( // 字符串去重后拼接 array_join(array_distinct(collect_list("company_name")), "; ").alias("company_name"), // 固定字段取首个值 first("industry").alias("industry"), // 数值保留所有唯一值 array_distinct(collect_list("revenue")).alias("revenue"), first("establish_date").alias("establish_date") ) mergedDF.show(false)
自定义UDF扩展版(针对复杂场景)
// 自定义字符串聚合UDF val concatUniqueUDF = udf((values: Seq[String]) => values.distinct.mkString("; ")) // 自定义数值聚合UDF(单个值返回原值,多值返回列表) val collectUniqueNumUDF = udf((values: Seq[Long]) => { val uniqueVals = values.distinct if (uniqueVals.size == 1) uniqueVals.head else uniqueVals }) val mergedDF = df.groupBy("company_id") .agg( concatUniqueUDF(collect_list("company_name")).alias("company_name"), first("industry").alias("industry"), collectUniqueNumUDF(collect_list("revenue")).alias("revenue"), first("establish_date").alias("establish_date") ) mergedDF.show(false)
优化建议
- 语义相似文本处理:若存在"ABC Corp"与"ABC Corporation"这类语义一致但字符串不同的情况,可在预处理阶段加入模糊匹配或同义词映射,统一值后再聚合
- 冲突字段优先级:同一主键下字段值冲突时,可定义规则(如最新数据优先、来源可靠数据优先),通过
row_number()标记优先级后筛选 - 性能优化:大数据量下,Spark先执行
df.dropDuplicates()过滤完全重复行再分组;Pandas同理先df.drop_duplicates()减少分组计算量
内容的提问来源于stack exchange,提问作者Georg Heiler
相关产品推荐
相关产品推荐

