基于Apache Beam实现Inner Join时对CoGroupByKey逻辑的疑问
Apache Beam Inner Join实现逻辑讲解
一、CoGroupByKey代码段作用解释
你提到的疑问代码如下:
PCollection<KV<String, CoGbkResult>> result = KeyedPCollectionTuple .of(customerOrderrTuple, customerOrderCollection).and(customerTuple, customerCollection) .apply(CoGroupByKey.<String>create());
分布式场景下两个数据集分散存储在不同计算节点,要实现按key关联的前提是把两个数据集中key相同的所有数据都汇聚到同一个计算节点,这段代码就是做这个核心预处理:
KeyedPCollectionTuple的作用是将你提前解析好的两个KV结构数据集(订单、用户)和对应的TupleTag绑定,相当于给两个数据集打唯一标识,避免后续聚合后分不清数据来源。CoGroupByKey执行按key聚合操作:同一个key的所有数据,不管来自订单数据集还是用户数据集,都会被打包到同一个CoGbkResult对象中,最终输出的KV<String, CoGbkResult>里,key是关联用的主键(一般是用户ID),value是该主键对应的两个数据集的全量数据集合。
这一步是所有分布式Join的必要前置操作,没有这一步跨数据集的key匹配根本无法完成。
二、ParDo处理段运行逻辑解释
你提到的疑问代码如下:
PCollection<String> output = result.apply(ParDo.of( new DoFn<KV<String, CoGbkResult>, String>() { @ProcessElement public void processElement(ProcessContext context) { String strKey = context.element().getKey(); CoGbkResult valueObject = context.element().getValue(); Iterable<String> customerOrderTable = valueObject.getAll(customerOrderrTuple); Iterable<String> customerTable = valueObject.getAll(customerTuple); for (String order:customerOrderTable) { for (String user:customerTable) { context.output(strKey+","+order+","+user); } } } } ));
这段代码是对聚合后的结果做遍历,实际生成Inner Join的最终输出,逻辑拆解:
- 逐个处理
CoGroupByKey输出的每一个KV元素:首先取出关联主键strKey,再取出该主键对应的聚合数据包CoGbkResult - 通过之前绑定的
TupleTag从聚合数据包中分别取出该主键对应的所有订单记录、所有用户记录 - 两层for循环对同一个key下的订单和用户记录做笛卡尔积拼接:
- 如果该主键只在订单数据集存在,那么
customerTable是空集合,内层循环不会执行,无输出 - 如果该主键只在用户数据集存在,那么
customerOrderTable是空集合,外层循环不会执行,无输出 - 只有该主键在两个数据集都存在时,才会将每一条订单和每一条用户记录做拼接输出,刚好符合Inner Join的语义。
- 如果该主键只在订单数据集存在,那么
内容的提问来源于stack exchange,提问作者User27854
相关产品推荐
相关产品推荐

