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

如何在Databricks中对关联物料分组并聚合库存(保留原有组名)

Databricks 物料关联分组与聚合解决方案

问题分析

  • 物料存在直接/间接关联关系,需将所有连通的物料归为同一组
  • 保留组内已有的非空组名称(如示例中的A1)
  • 聚合同组物料的STOCKS总量
  • 最终结果存储至ADLS

模拟示例数据结构:

MATERIALRELATED_MATERIALGROUP_NAMESTOCKS
ABA110
BCNULL20
CNULLNULL15
DED15
ENULLNULL8

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 21:37:26