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

ORBIS企业多重复键值数据的通用整合策略咨询

ORBIS企业数据主键多行合并通用策略(附Pandas/Spark 3.5实现)

场景与需求

处理ORBIS企业数据时,存在同一主键(如company_id)对应多行数据的情况,包含三类数据:

  • 完全重复的行(键值+字段内容均一致)
  • 同一主键下语义相似但字段有差异的行(需保留多值)
  • 同主键补充字段或异键的新数据

目标是将同一主键的所有有效数据合并为单行,既要去除冗余重复,又要保留同一字段的多值信息(避免仅保留首个值的局限)。

通用整合策略

  1. 预处理去重:先过滤完全重复的行,减少后续计算量
  2. 字段类型适配聚合:针对不同字段类型定义专属聚合规则:
    • 字符串/枚举型:收集唯一值后用分隔符拼接,或保留所有有效语义值
    • 数值型:根据业务需求选择保留所有唯一值、求和、均值或最值
    • 固定属性字段(如行业、成立日期):取首个值或统一值
  3. 自定义逻辑扩展:针对语义相似的文本字段,可加入模糊匹配、同义词映射等预处理逻辑,再进行聚合

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 23:27:20