Flink主流数据连接广播流数据的水位正常生成处理
·
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

更多推荐


所有评论(0)