Flink主流数据连接广播流数据的水位正常生成处理

当使用Flink的主流数据来连接广播流数据的时候 会因为维度流数据不包含水位线 导致connect后无法正常生成水位线 进而导致无法触发窗口计算或定时器 因为水位线的传递策略就是这样 下游的水位线依靠上游所有流的水位的最小值

因此只要在维度流中指定水位线的生成策略即可保证connect后也能够正常生成 考虑到水位线取最小值这个特性 可以给广播流数据恒定一个无法达到的水位线 此时连接后的水位线策略就完全取决于业务流数据的时间

因此就有以下内容

package Tips;

import lombok.AllArgsConstructor;
import lombok.Data;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.BroadcastStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction;
import org.apache.flink.util.Collector;

import java.time.Duration;

/**
 * @Auther mengyi
 * @Date 2025/10/9 22:36
 * Project Demo flink
 */
public class ConnectBroadcastWatermarkDemo {
    @AllArgsConstructor
    @Data
    public static class Dim{
        private Integer id;
        private String name;
    }
    @AllArgsConstructor
    @Data
    public static class Business{
        private Integer id;
        private Long timeStamp;
    }
    public static class String2Business extends ProcessFunction<String,Business>{
        @Override
        public void processElement(String value, ProcessFunction<String, Business>.Context ctx, Collector<Business> out) throws Exception {
            String [] fields = value.split(",");
            if(fields.length==2){
                out.collect(new Business(Integer.parseInt(fields[0]),Long.parseLong(fields[1])));
            }
        }
    }
    public static class String2Dim extends ProcessFunction<String,Dim>{
        @Override
        public void processElement(String value, ProcessFunction<String, Dim>.Context ctx, Collector<Dim> out) throws Exception {
            String[] fields = value.split(",");
            if(fields.length==2){
                out.collect(new Dim(Integer.parseInt(fields[0]),fields[1]));
            }
        }
    }
    public static class BusinessJoinDimProcessFunction extends BroadcastProcessFunction<Business,Dim,String> {
        @Override
        public void processElement(Business value, BroadcastProcessFunction<Business, Dim, String>.ReadOnlyContext ctx, Collector<String> out) throws Exception {
            System.out.println("ctx.currentWatermark() = " + ctx.currentWatermark());
        }

        @Override
        public void processBroadcastElement(Dim value, BroadcastProcessFunction<Business, Dim, String>.Context ctx, Collector<String> out) throws Exception {
            System.out.println("ctx.currentWatermark() = " + ctx.currentWatermark());
        }
    }
    public static void main(String[] args) throws Exception {
        final MapStateDescriptor<String, Dim> dimMapStateDescriptor = new MapStateDescriptor<>("dimMapStateDesc", String.class,Dim.class);
        Configuration conf = new Configuration();
        conf.setInteger("rest.port",2516);
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);
        env.setParallelism(1).enableCheckpointing(60000L);
        BroadcastStream<Dim> dimBroadcastStream = env
                .socketTextStream("localhost", 1235)
                .process(new String2Dim())
                .assignTimestampsAndWatermarks(WatermarkStrategy.<Dim>forBoundedOutOfOrderness(Duration.ofMillis(1))
                .withTimestampAssigner((SerializableTimestampAssigner<Dim>) (element, recordTimestamp) -> 2999999999999L))
                .broadcast(dimMapStateDescriptor);
//      读取主流连接广播的维度数据
        env
            .socketTextStream("localhost",1234)
            .process(new String2Business())
            .assignTimestampsAndWatermarks(WatermarkStrategy.<Business>forBoundedOutOfOrderness(Duration.ofMillis(1))
            .withTimestampAssigner((SerializableTimestampAssigner<Business>) (element, recordTimestamp) -> element.getTimeStamp()))
            .connect(dimBroadcastStream)
            .process(new BusinessJoinDimProcessFunction());
        env.execute();
    }
}
-- 控制台开启两个端口
> nc -Lp 1235
1,1
>nc -Lp 1234
1,1
1,2
1,3
得到数据如下
ctx.currentWatermark() = -9223372036854775808
ctx.currentWatermark() = -9223372036854775808
ctx.currentWatermark() = -1
ctx.currentWatermark() = 0

在这里插入图片描述

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐