Flink批模式DataStream API下DataSet API Left Join的等效实现方法
Flink DataStream 批模式左外连接实现方案
Flink DataStream 批模式下支持两种常用的左外连接实现方式,均可对齐原有 DataSet API 的 leftOuterJoin 语义:
方案1:通过 Table/SQL API 中转实现(推荐)
批模式下 Table API 与 DataStream 可无缝互转,代码实现简洁,底层自动适配批执行优化,也支持自定义Join策略:
- 首先将两个 DataStream 转为 Table 对象,声明所需业务字段
- 编写左外连接 SQL 完成关联,关联条件可直接复用原有的
coalesce(left.getId(), -9999999L) = right.company_id逻辑,也可通过查询Hint指定广播策略对齐原有的BROADCAST_HASH_SECOND配置 - 最后将关联结果 Table 转回指定类型的 DataStream 即可
示例代码如下:
// 初始化批模式适配的Table环境 StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env); // DataStream转Table,声明所需字段 Table tableA = tableEnv.fromDataStream(datasetA, $("id"), $("customerId")); Table tableB = tableEnv.fromDataStream(datasetB, $("company_id"), $("cust_name")); // 执行带广播Hint的左外连接 Table joinedTable = tableEnv.sqlQuery( "SELECT /*+ BROADCAST(b) */ a.customerId, COALESCE(b.cust_name, 'Blank') as cust_name " + "FROM " + tableA + " a LEFT JOIN " + tableB + " b " + "ON COALESCE(a.id, -9999999L) = b.company_id" ); // 关联结果转回DataStream DataStream<SomeOutType> joinedOut = tableEnv.toDataStream(joinedTable, SomeOutType.class);
方案2:纯DataStream API 原生实现
如果不想引入Table/SQL依赖,可通过自定义KeyedCoProcessFunction实现左外连接逻辑:
- 分别对两个DataStream按关联键做keyBy
- 用connect算子连接两个KeyedStream
- 自定义KeyedCoProcessFunction,用状态缓存右流的同key数据,左流每来一条数据就匹配缓存的右流数据输出,无匹配则输出右值为null的结果
- 批模式下无需处理流场景的窗口、迟到数据等问题,状态会在作业结束后自动清理
示例代码如下:
// 按关联键分别做keyBy KeyedStream<SomeTypeA, Long> keyedA = datasetA .keyBy(left -> Objects.nonNull(left.getId()) ? left.getId() : -9999999L); KeyedStream<SomeTypeB, Long> keyedB = datasetB .keyBy(SomeTypeB::getCompany_id); DataStream<SomeOutType> joinedOut = keyedA.connect(keyedB) .process(new KeyedCoProcessFunction<Long, SomeTypeA, SomeTypeB, SomeOutType>() { private ListState<SomeTypeB> bDataState; @Override public void open(Configuration parameters) throws Exception { bDataState = getRuntimeContext().getListState( new ListStateDescriptor<>("b_cache", SomeTypeB.class) ); } @Override public void processElement1(SomeTypeA left, Context ctx, Collector<SomeOutType> out) throws Exception { boolean hasMatch = false; // 遍历所有匹配的右流数据输出 for (SomeTypeB right : bDataState.get()) { hasMatch = true; out.collect(buildOutput(left, right)); } // 无匹配则输出右值为null的结果 if (!hasMatch) { out.collect(buildOutput(left, null)); } } @Override public void processElement2(SomeTypeB right, Context ctx, Collector<SomeOutType> out) throws Exception { // 右流数据存入状态缓存 bDataState.add(right); } private SomeOutType buildOutput(SomeTypeA left, SomeTypeB right) { SomeOutType output = SomeOutType.newBuilder().build(); output.setCustomerId(left.getCustomerId()); output.setCustomerName(Objects.nonNull(right) && Objects.nonNull(right.getCust_name()) ? right.getCust_name() : "Blank"); // 其他字段赋值逻辑 return output; } });
注意:纯DataStream实现如果右流数据量较大,状态内存占用会较高,优先推荐使用Table/SQL方案,底层优化更成熟,性能与原DataSet API接近。
内容的提问来源于stack exchange,提问作者ybd
相关产品推荐
相关产品推荐

