如何在Databricks中对关联物料分组并聚合库存(保留原有组名)
Databricks 物料关联分组与聚合解决方案
问题分析
- 物料存在直接/间接关联关系,需将所有连通的物料归为同一组
- 保留组内已有的非空组名称(如示例中的A1)
- 聚合同组物料的STOCKS总量
- 最终结果存储至ADLS
模拟示例数据结构:
MATERIAL RELATED_MATERIAL GROUP_NAME STOCKS A B A1 10 B C NULL 20 C NULL NULL 15 D E D1 5 E NULL NULL 8
SQL 实现方案
1. 加载数据并创建临时视图
-- 替换为实际数据存储路径(如ADLS上的Delta表路径) CREATE OR REPLACE TEMP VIEW material_data AS SELECT * FROM delta.`/path/to/your/material/source`;
2. 递归CTE处理连通分量与分组
WITH RECURSIVE material_connections AS ( -- 初始化:每个物料的根节点设为自身 SELECT MATERIAL AS node, MATERIAL AS root, GROUP_NAME FROM material_data UNION ALL -- 递归关联:通过RELATED_MATERIAL追溯根节点 SELECT mc.node, md.MATERIAL AS root, mc.GROUP_NAME FROM material_connections mc JOIN material_data md ON mc.root = md.RELATED_MATERIAL WHERE mc.root != md.MATERIAL -- 避免循环引用 ), -- 为每个连通分量确定最终组名称(优先取非空值) group_name_mapping AS ( SELECT node, FIRST_VALUE(GROUP_NAME) OVER ( PARTITION BY root ORDER BY CASE WHEN GROUP_NAME IS NOT NULL THEN 0 ELSE 1 END ) AS final_group_name FROM material_connections ) -- 聚合同组物料的库存 SELECT gnm.final_group_name AS GROUP_NAME, SUM(md.STOCKS) AS TOTAL_STOCKS, COLLECT_SET(md.MATERIAL) AS GROUP_MATERIALS FROM group_name_mapping gnm JOIN material_data md ON gnm.node = md.MATERIAL GROUP BY gnm.final_group_name;
3. 将结果写入ADLS
-- 替换为你的ADLS存储路径(ABFS格式) INSERT INTO delta.`abfss://container@your-storage-account.dfs.core.windows.net/path/to/output` SELECT gnm.final_group_name AS GROUP_NAME, SUM(md.STOCKS) AS TOTAL_STOCKS, COLLECT_SET(md.MATERIAL) AS GROUP_MATERIALS FROM group_name_mapping gnm JOIN material_data md ON gnm.node = md.MATERIAL GROUP BY gnm.final_group_name;
Python 实现方案
1. 加载数据到DataFrame
from pyspark.sql import functions as F from pyspark.sql.window import Window -- 替换为实际数据路径 df = spark.read.format("delta").load("/path/to/your/material/source")
2. 递归计算连通分量的根节点
def get_material_roots(raw_df, max_iter=15): # 初始化:每个物料的根节点为自身 processed_df = raw_df.withColumn("root", F.col("MATERIAL")) for _ in range(max_iter): # 关联父节点,更新根节点 processed_df = processed_df.alias("curr").join( processed_df.alias("parent"), F.col("curr.RELATED_MATERIAL") == F.col("parent.MATERIAL"), "left" ).select( F.col("curr.MATERIAL"), F.col("curr.RELATED_MATERIAL"), F.col("curr.GROUP_NAME"), F.col("curr.STOCKS"), F.coalesce(F.col("parent.root"), F.col("curr.root")).alias("root") ) return processed_df root_df = get_material_roots(df)
3. 确定组名称并聚合库存
-- 窗口函数:按根节点分组,优先取非空的组名称 window_spec = Window.partitionBy("root").orderBy(F.col("GROUP_NAME").isNull().cast("int")) grouped_result = root_df.withColumn( "final_group_name", F.first(F.col("GROUP_NAME")).over(window_spec) ).groupBy("final_group_name").agg( F.sum("STOCKS").alias("TOTAL_STOCKS"), F.collect_set("MATERIAL").alias("GROUP_MATERIALS") ) -- 查看结果 grouped_result.show()
4. 写入ADLS
-- 替换为你的ADLS路径,支持Delta/Parquet等格式 grouped_result.write.format("delta").mode("overwrite").save( "abfss://container@your-storage-account.dfs.core.windows.net/path/to/output" )
注意事项
- 若物料关联层级较深,可调整递归迭代次数(Python方案)或确保SQL递归CTE无循环
- 推荐使用Delta Lake格式存储结果,支持ACID事务与版本管理
- 需确保Databricks集群已配置ADLS访问权限(服务主体/SAS令牌)
内容的提问来源于stack exchange,提问作者Eleazzz
相关产品推荐
相关产品推荐

