You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink批模式DataStream API下DataSet API Left Join的等效实现方法

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 12:39:03