如何通过Dataflow Pipeline的PCollection关联BigQuery两张表获取数据?
问题描述
我有两个BigQuery表,结构如下:
表A:
c_id count_c_id p_id
表B:
id c_name p_type c_id
我已经用Dataflow Pipeline把表A的数据读取到PCollection<TableRow>里了,代码是:
PCollection<TableRow> tableRowBQ = pipeline.apply(BigQueryIO.Read.named("Read").fromQuery("select c_id,count_c_id,p_id from TableA"));
现在我想利用表A里的c_id,从表B中获取对应的c_name,但找不到怎么通过PCollection迭代字段去另一张表取数的示例,参考过官方的Join示例还是没理清具体实现方法。
解决方案
其实你要实现的就是Dataflow里的关联操作,这里给你两种可行的方案,按需选择:
方案一:在Dataflow中用侧输入(Side Input)关联
这种方式适合表B数据量不大的场景,把表B作为维度表加载成可查询的侧输入,再和表A的主数据匹配:
步骤1:加载表B并转换成侧输入Map
先把表B的数据转换成KV<c_id, c_name>的格式,再生成一个可被查询的Map侧输入:
// 读取表B,提取c_id和c_name并转成KV对 PCollection<KV<String, String>> bTableKV = pipeline.apply( BigQueryIO.Read.named("ReadTableB") .fromQuery("select c_id, c_name from TableB")) .apply(ParDo.of(new DoFn<TableRow, KV<String, String>>() { @Override public void processElement(ProcessContext c) { TableRow row = c.element(); String cId = row.get("c_id").toString(); String cName = row.get("c_name").toString(); c.output(KV.of(cId, cName)); } })); // 将KV集合转换为侧输入Map,方便后续快速查找 final PCollectionView<Map<String, String>> cIdToNameMap = bTableKV.apply(View.asMap());
步骤2:处理表A数据并关联侧输入
在处理表A的PCollection时,传入侧输入的Map,通过c_id直接匹配对应的c_name:
PCollection<TableRow> joinedResults = tableRowBQ.apply( ParDo.of(new DoFn<TableRow, TableRow>() { @Override public void processElement(ProcessContext c) { TableRow aRow = c.element(); String cId = aRow.get("c_id").toString(); // 从侧输入Map中获取对应c_name,处理找不到的情况 String cName = c.sideInput(cIdToNameMap).getOrDefault(cId, "N/A"); // 构造包含c_name的结果行 TableRow resultRow = new TableRow(); resultRow.set("c_id", cId); resultRow.set("count_c_id", aRow.get("count_c_id")); resultRow.set("p_id", aRow.get("p_id")); resultRow.set("c_name", cName); c.output(resultRow); } }).withSideInputs(cIdToNameMap));
方案二:直接在BigQuery层面做JOIN查询
如果表的数据量比较大,更推荐直接在BigQuery里完成关联,因为BigQuery本身对大数据关联的优化更好,性能更高,还能减少Dataflow的处理压力:
// 直接读取BigQuery关联后的结果,省去Dataflow中的关联步骤 PCollection<TableRow> joinedBQ = pipeline.apply( BigQueryIO.Read.named("ReadJoinedTables") .fromQuery("SELECT a.c_id, a.count_c_id, a.p_id, b.c_name " + "FROM TableA a LEFT JOIN TableB b ON a.c_id = b.c_id"));
额外提示
- 如果表B数据量极大,侧输入的Map可能会占用过多内存,这时候可以考虑使用
CoGroupByKey做关联,但需要把两个表都转换成KV格式后按c_id分组,不过这种方式的复杂度比侧输入高一些。 - 用BigQuery JOIN的方式更简洁高效,除非你有特殊需求必须在Dataflow里做数据处理,否则优先选这个方案。
内容的提问来源于stack exchange,提问作者user1115163
相关产品推荐
相关产品推荐

