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

如何使用Azure Data Factory汇总动态字段并聚合库存数据

问题描述

现有如下结构的表格:

id, stockplace, stock
1, A-100, 40
1, B-100, 10
1, C-100, 5
2, A-300, 90

表格包含商品ID、库存地点及库存数量,部分商品会分布在多个库存地点且库存数量不同。需要生成最终表格,按唯一ID汇总总库存,并合并该商品对应的所有库存地点,示例如下:

id, stockplace, stock
1, A-100 & B-100 & C-100, 55
2, A-300, 90

此前尝试用窗口流按id分组,添加rowNumber()窗口列,用sum(stock)表达式计算唯一ID的总库存,但未成功,希望得到解决思路或指导。

解决思路

1. 传统SQL场景

直接按id分组,借助聚合函数即可实现需求:

  • 总库存:用SUM(stock)直接求和
  • 库存地点合并:用字符串聚合函数,不同数据库对应函数不同:
    • MySQL 5.7+:GROUP_CONCAT(stockplace SEPARATOR ' & ')
    • PostgreSQL/SQL Server 2017+:STRING_AGG(stockplace, ' & ')

示例SQL(以MySQL为例):

SELECT 
  id,
  GROUP_CONCAT(stockplace SEPARATOR ' & ') AS stockplace,
  SUM(stock) AS stock
FROM your_table_name
GROUP BY id;

2. 流处理场景(以Flink为例)

如果是用Flink这类流处理框架,无需使用rowNumber(),直接基于id分组做聚合即可:

  • 总库存:用内置的SUM("stock")聚合
  • 库存地点合并:用字符串聚合函数或自定义聚合逻辑
-- 定义输入表
CREATE TABLE input_stock (
  id INT,
  stockplace STRING,
  stock INT
) WITH (
  -- 数据源配置,比如Kafka、文件等
);

-- 定义输出表
CREATE TABLE output_stock (
  id INT,
  stockplace STRING,
  stock INT
) WITH (
  -- 输出配置
);

-- 聚合计算并写入输出表
INSERT INTO output_stock
SELECT 
  id,
  STRING_AGG(stockplace, ' & ') AS stockplace,
  SUM(stock) AS stock
FROM input_stock
GROUP BY id;

如果用DataStream API,可自定义聚合函数实现:

// 定义输入实体类
public class StockRecord {
    private Integer id;
    private String stockplace;
    private Integer stock;
    // 构造器、getter、setter省略
}

// 定义累加器
public class StockAccumulator {
    public Integer totalStock;
    public List<String> places;

    public StockAccumulator(Integer totalStock, List<String> places) {
        this.totalStock = totalStock;
        this.places = places;
    }
}

// 定义输出结果类
public class StockResult {
    public Integer id;
    public String stockplace;
    public Integer stock;

    public StockResult(Integer id, String stockplace, Integer stock) {
        this.id = id;
        this.stockplace = stockplace;
        this.stock = stock;
    }
}

// 自定义聚合函数
public class StockAggregate implements AggregateFunction<StockRecord, StockAccumulator, StockResult> {
    @Override
    public StockAccumulator createAccumulator() {
        return new StockAccumulator(0, new ArrayList<>());
    }

    @Override
    public void add(StockRecord value, StockAccumulator accumulator) {
        accumulator.totalStock += value.getStock();
        accumulator.places.add(value.getStockplace());
    }

    @Override
    public StockResult getResult(StockAccumulator accumulator) {
        String placeStr = String.join(" & ", accumulator.places);
        return new StockResult(null, placeStr, accumulator.totalStock);
    }

    @Override
    public StockAccumulator merge(StockAccumulator a, StockAccumulator b) {
        a.totalStock += b.totalStock;
        a.places.addAll(b.places);
        return a;
    }
}

// 处理逻辑
DataStream<StockRecord> inputStream = ...; // 读取数据源
inputStream.keyBy(StockRecord::getId)
           .aggregate(new StockAggregate())
           .map(result -> {
               // 补充id,可从keyBy上下文获取或调整聚合逻辑
               return new StockResult(...);
           })
           .addSink(...); // 输出到目标

失败原因分析

之前用窗口流加rowNumber()的思路有误,rowNumber()是用于给窗口内的行排序编号,而需求是分组聚合,无需行号列。直接按id分组后执行聚合操作,就能得到目标结果。

内容的提问来源于stack exchange,提问作者Julian

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:48:40