如何在Pentaho DI中不阻塞流添加非常量新列并实现笛卡尔积效果
在Pentaho DI中实现非阻塞流的笛卡尔积连接
核心思路
无需等待全量数据流加载完成,通过缓存次级流数据,让主流逐行匹配缓存内的所有次级行,或借助行列转换逻辑,实现无阻塞的笛卡尔积效果。
具体实现方案(非阻塞版本)
方案1:流查询+内存缓存的循环匹配(推荐)
全程不阻塞主流,先将次级流缓存至内存,主流逐行循环匹配所有次级行:
- 处理次级节点数据
- 用
Memory Group by组件将次级节点的所有输出行缓存到内存(若次级数据量过大,可替换为Cache组件使用磁盘缓存,内存缓存性能更优) - 给次级流每一行添加唯一序号,通过
Add sequence组件生成seq_id,从1开始自增
- 用
- 处理主流节点数据
- 给主流每一行添加匹配序号
match_seq,初始值设为1 - 提前用
Row count组件统计次级数据总行数,存入全局变量SUB_ROW_COUNT,再添加Loop组件,设置循环条件为match_seq <= ${SUB_ROW_COUNT}
- 给主流每一行添加匹配序号
- 循环内完成笛卡尔匹配
- 循环内部使用
Stream Lookup组件,以match_seq = seq_id为匹配条件,从缓存的次级流中取出对应行,与当前主流行合并 - 合并后的行直接输出结果,再通过
Calculator组件将match_seq自增1,回到循环入口 - 当
match_seq超过次级总行数时,结束循环,主流行进入后续处理环节
- 循环内部使用
方案2:行转列+列转行的无阻塞实现
若次级流字段较少,可通过该方案省去循环逻辑:
- 处理次级流
- 用
Pivot rows组件将次级流的所有行转换为一行多列格式,例如把次级的col1值用分隔符(如|)拼接成col1_1|col1_2|...|col1_n,col2同理转换为col2_1|col2_2|...|col2_n - 将转换后的单行数据通过
Set variable存入全局变量,或用Copy rows to result传递给主流
- 用
- 处理主流流
- 给主流添加
Split fields to rows组件,将全局变量中的次级多列数据按分隔符拆分为多行,每一行对应次级流的一条原始数据,从而与主流行形成笛卡尔积 - 拆分时需保证所有次级字段的拆分规则一致,统一使用相同分隔符
- 给主流添加
关键注意事项
- 若次级数据量极大,避免使用内存缓存,改用磁盘缓存防止内存溢出
- 次级行总数需提前统计完成,确保全局变量在主流处理前已设置完毕
- 两种方案均为逐行处理主流,不会等待全量主流数据加载,完全满足非阻塞要求
内容的提问来源于stack exchange,提问作者Imangali
相关产品推荐
相关产品推荐

