如何使用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, ' & ')
- MySQL 5.7+:
示例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")聚合 - 库存地点合并:用字符串聚合函数或自定义聚合逻辑
Flink SQL示例
-- 定义输入表 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;
Flink DataStream API示例
如果用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
相关产品推荐
相关产品推荐

