Flink经典编程场景实战解析
1: 分组TOPN
代码思路
1. 定义单独的POJO类 UserBehavior 和 ItemViewCount
UserBehavior:解析JSON字符串后生成的JavaBean。ItemViewCount:最终结果输出的格式类。
2. 调用底层的 Process 方法
- 将JSON字符串解析成
UserBehavior对象。
3. 提取 EventTime 并转换为 Timestamp 格式,生成 WaterMark。
4. 按照指定事件分组。
5. 划分窗口
- 假设窗口总长为10分钟,步长为1分钟滑动一次。
6. 调用 aggregate 方法,在窗口内增量聚合
MyWindowAggFunction:获取聚合字段(UserBehavior中的counts)。- 三个泛型:输入类型、累加器类型、输出数据类型。
MyWindowFunction:获取窗口的开始时间和结束时间,提取分组字段。- 四个泛型:输入数据类型(
Long类型的次数)、输出数据类型(ItemViewCount)、分组的key、窗口对象。
- 四个泛型:输入数据类型(
7. 对聚合好的窗口内数据排序
- 按照窗口的
start和end进行分组,将窗口相同的数据进行排序。
代码示例
java
import lombok.Data;
@Data
public class UserBehavior {
public String userId; // 用户ID
public String itemId; // 商品ID
public String categoryId; // 商品类目ID
public String type; // 用户行为, 包括("pv", "buy", "cart", "fav")
public long timestamp; // 行为发生的时间戳,单位秒
public long counts = 1;
public static UserBehavior of(String userId, String itemId, String categoryId, String type, long timestamp) {
UserBehavior behavior = new UserBehavior();
behavior.userId = userId;
behavior.itemId = itemId;
behavior.categoryId = categoryId;
behavior.type = type;
behavior.timestamp = timestamp;
return behavior;
}
public static UserBehavior of(String userId, String itemId, String categoryId, String type, long timestamp, long counts) {
UserBehavior behavior = new UserBehavior();
behavior.userId = userId;
behavior.itemId = itemId;
behavior.categoryId = categoryId;
behavior.type = type;
behavior.timestamp = timestamp;
behavior.counts = counts;
return behavior;
}
}
@Data
public class ItemViewCount {
public String itemId; // 商品ID
public String type; // 事件类型
public long windowStart; // 窗口开始时间戳
public long windowEnd; // 窗口结束时间戳
public long viewCount; // 商品的点击量
public static ItemViewCount of(String itemId, String type, long windowStart, long windowEnd, long viewCount) {
ItemViewCount result = new ItemViewCount();
result.itemId = itemId;
result.type = type;
result.windowStart = windowStart;
result.windowEnd = windowEnd;
result.viewCount = viewCount;
return result;
}
}
java
import com.chehejia.dip.pojo.ItemViewCount;
import org.apache.flink.api.java.tuple.Tuple;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
public class MyWindowFunction implements WindowFunction<Long, ItemViewCount, Tuple, TimeWindow> {
@Override
public void apply(Tuple tuple, TimeWindow window, Iterable<Long> input, Collector<ItemViewCount> out) throws Exception {
String itemId = tuple.getField(0);
String type = tuple.getField(1);
long windowStart = window.getStart();
long windowEnd = window.getEnd();
Long aLong = input.iterator().next();
out.collect(ItemViewCount.of(itemId, type, windowStart, windowEnd, aLong));
}
}
java
import com.chehejia.dip.pojo.UserBehavior;
import org.apache.flink.api.common.functions.AggregateFunction;
public class MyWindowAggFunction implements AggregateFunction<UserBehavior, Long, Long> {
@Override
public Long createAccumulator() {
return 0L;
}
@Override
public Long add(UserBehavior input, Long accumulator) {
return accumulator + input.counts;
}
@Override
public Long getResult(Long accumulator) {
return accumulator;
}
@Override
public Long merge(Long a, Long b) {
return null;
}
}
java
import com.alibaba.fastjson.JSON;
import com.chehejia.dip.pojo.ItemViewCount;
import com.chehejia.dip.pojo.UserBehavior;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.KeyedStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.datastream.WindowedStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor;
import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
public class HotGoodsTopN {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env.enableCheckpointing(60000);
env.setParallelism(1);
DataStreamSource<String> lines = env.socketTextStream("linux01", 8888);
SingleOutputStreamOperator<UserBehavior> process = lines.process(new ProcessFunction<String, UserBehavior>() {
@Override
public void processElement(String input, Context ctx, Collector<UserBehavior> out) throws Exception {
try {
UserBehavior behavior = JSON.parseObject(input, UserBehavior.class);
out.collect(behavior);
} catch (Exception e) {
e.printStackTrace();
}
}
});
SingleOutputStreamOperator<UserBehavior> behaviorDSWithWaterMark = process.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<UserBehavior>(Time.seconds(0)) {
@Override
public long extractTimestamp(UserBehavior element) {
return element.timestamp;
}
});
KeyedStream<UserBehavior, Tuple> keyed = behaviorDSWithWaterMark.keyBy("itemId", "type");
WindowedStream<UserBehavior, Tuple, TimeWindow> window = keyed.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)));
SingleOutputStreamOperator<ItemViewCount> windowAggregate = window.aggregate(new MyWindowAggFunction(), new MyWindowFunction());
KeyedStream<ItemViewCount, Tuple> soredKeyed = windowAggregate.keyBy("type", "windowStart", "windowEnd");
SingleOutputStreamOperator<List<ItemViewCount>> sored = soredKeyed.process(new KeyedProcessFunction<Tuple, ItemViewCount, List<ItemViewCount>>() {
private transient ValueState<List<ItemViewCount>> valueState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<List<ItemViewCount>> VSDescriptor = new ValueStateDescriptor<>("list-state", TypeInformation.of(new TypeHint<List<ItemViewCount>>() {}));
valueState = getRuntimeContext().getState(VSDescriptor);
}
@Override
public void processElement(ItemViewCount input, Context ctx, Collector<List<ItemViewCount>> out) throws Exception {
List<ItemViewCount> buffer = valueState.value();
if (buffer == null) {
buffer = new ArrayList<>();
}
buffer.add(input);
valueState.update(buffer);
ctx.timerService().registerEventTimeTimer(input.windowEnd + 1);
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<List<ItemViewCount>> out) throws Exception {
List<ItemViewCount> buffer = valueState.value();
buffer.sort(new Comparator<ItemViewCount>() {
@Override
public int compare(ItemViewCount o1, ItemViewCount o2) {
return -(int) (o1.viewCount - o2.viewCount);
}
});
valueState.update(null);
out.collect(buffer);
}
});
env.execute("HotGoodsTopNAdv");
}
}
2: Flink使用二次聚合实现TopN计算
需求背景
需要每隔5分钟输出最近1小时内点击量最多的前N个商品。
实现思路
- 建立环境,设置并行度及Checkpoint。
- 定义Watermark策略及事件时间,获取数据并对应到JavaBean,筛选PV数据。
- 第一次聚合,按商品ID分组开窗聚合,使用
aggregate算子进行增量计算。 - 第二次聚合,按窗口聚合,使用
ListState存放数据,并定义定时器,在Watermark达到后1秒触发,对窗口数据排序输出。 - 打印结果并执行。
代码示例
java
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor
@AllArgsConstructor
public class ApacheLog {
private String ip;
private String userId;
private Long ts;
private String method;
private String url;
}
@Data
@NoArgsConstructor
@AllArgsConstructor
public class UrlCount {
private String url;
private Long windowEnd;
private Integer count;
}
java
package com.test.topN;
import bean.ApacheLog;
import bean.UrlCount;
import org.apache.commons.compress.utils.Lists;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.AggregateFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import java.sql.Timestamp;
import java.text.SimpleDateFormat;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.Map;
public class URLTopN3 {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment().setParallelism(1);
WatermarkStrategy<ApacheLog> wms = WatermarkStrategy.<ApacheLog>forBoundedOutOfOrderness(Duration.ofSeconds(1))
.withTimestampAssigner(new SerializableTimestampAssigner<ApacheLog>() {
@Override
public long extractTimestamp(ApacheLog element, long recordTimestamp) {
return element.getTs();
}
});
SingleOutputStreamOperator<ApacheLog> apacheLogDS = env.socketTextStream("hadoop102", 9999)
.map(new MapFunction<String, ApacheLog>() {
@Override
public ApacheLog map(String value) throws Exception {
SimpleDateFormat sdf = new SimpleDateFormat("dd/MM/yy:HH:mm:ss");
String[] split = value.split(" ");
return new ApacheLog(split[0], split[2], sdf.parse(split[3]).getTime(), split[5], split[6]);
}
})
.assignTimestampsAndWatermarks(wms);
SingleOutputStreamOperator<UrlCount> aggregateDS = apacheLogDS
.map(new MapFunction<ApacheLog, Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> map(ApacheLog value) throws Exception {
return new Tuple2<>(value.getUrl(), 1);
}
}).keyBy(data -> data.f0)
.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.seconds(5)))
.allowedLateness(Time.minutes(1))
.aggregate(new HotUrlAggFunc(), new HotUrlWindowFunc());
SingleOutputStreamOperator<String> processDS = aggregateDS
.keyBy(data -> data.getWindowEnd())
.process(new HotUrlProcessFunc(5));
processDS.print();
env.execute();
}
public static class HotUrlAggFunc implements AggregateFunction<Tuple2<String, Integer>, Integer, Integer> {
@Override
public Integer createAccumulator() {
return 0;
}
@Override
public Integer add(Tuple2<String, Integer> value, Integer accumulator) {
return accumulator + 1;
}
@Override
public Integer getResult(Integer accumulator) {
return accumulator;
}
@Override
public Integer merge(Integer a, Integer b) {
return a + b;
}
}
public static class HotUrlWindowFunc implements WindowFunction<Integer, UrlCount, String, TimeWindow> {
@Override
public void apply(String urls, TimeWindow window, Iterable<Integer> input, Collector<UrlCount> out) throws Exception {
Integer count = input.iterator().next();
out.collect(new UrlCount(urls, window.getEnd(), count));
}
}
public static class HotUrlProcessFunc extends KeyedProcessFunction<Long, UrlCount, String> {
private Integer TopN;
private MapState<String, UrlCount> mapState;
public HotUrlProcessFunc(Integer topN) {
TopN = topN;
}
@Override
public void open(Configuration parameters) throws Exception {
mapState = getRuntimeContext().getMapState(new MapStateDescriptor<String, UrlCount>("map-state", String.class, UrlCount.class));
}
@Override
public void processElement(UrlCount value, Context ctx, Collector<String> out) throws Exception {
mapState.put(value.getUrl(), value);
ctx.timerService().registerEventTimeTimer(value.getWindowEnd() + 1L);
ctx.timerService().registerEventTimeTimer(value.getWindowEnd() + 61001L);
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) throws Exception {
if (timestamp == ctx.getCurrentKey() + 61001L) {
mapState.clear();
return;
}
Iterator<Map.Entry<String, UrlCount>> iterator = mapState.iterator();
ArrayList<Map.Entry<String, UrlCount>> entries = Lists.newArrayList(iterator);
entries.sort(((o1, o2) -> o2.getValue().getCount() - o1.getValue().getCount()));
StringBuilder sb = new StringBuilder();
sb.append("==============")
.append(new Timestamp(timestamp - 1L))
.append("==============")
.append("
");
for (int i = 0; i < Math.min(TopN, entries.size()); i++) {
UrlCount urlCount = entries.get(i).getValue();
sb.append("Top").append(i + 1);
sb.append(" Url:").append(urlCount.getUrl());
sb.append(" Counts:").append(urlCount.getCount());
sb.append("
");
}
sb.append("==============")
.append(new Timestamp(timestamp - 1L))
.append("==============")
.append("
")
.append("
");
out.collect(sb.toString());
Thread.sleep(200);
}
}
}
3: PV、UV统计
需求描述
从Kafka发送过来的数据含有:时间戳、时间、维度、用户ID,需要从不同维度统计从0点到当前时间的PV和UV,第二天0点重新开始计数第二天的。
- PV(访问量):即Page View,即页面浏览量或点击量,用户每次刷新即被计算一次。
- UV(独立访客):即Unique Visitor,访问您网站的一台电脑客户端为一个访客。00:00-24:00内相同的客户端只被计算一次。
实现思路
- Kafka数据可能会有延迟乱序,这里引入Watermark。
- 通过
keyBy分流进不同的滚动窗口,每个窗口内计算PV、UV。 - 由于需要保存一天的状态,
process里面使用ValueState保存PV、UV。 - 使用
BitMap类型ValueState,占内存很小,引入支持bitmap的依赖。 - 保存状态需要设置TTL过期时间,第二天把第一天的过期,避免内存占用过大。
代码实现
java
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.ToString;
@Data
@NoArgsConstructor
@AllArgsConstructor
@ToString
public class UserClickModel {
private String date;
private String product;
private int uid;
private int pv;
private int uv;
}
java
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
import org.apache.flink.streaming.api.functions.source.SourceFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
public class UserClickMain {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env.enableCheckpointing(5000, CheckpointingMode.EXACTLY_ONCE);
env.setStateBackend(new FsStateBackend("hdfs://bigdata/flink/checkpoints/userClick"));
Properties kafkaProps = new Properties();
kafkaProps.setProperty("bootstrap.servers", config.get("kafka-ipport"));
kafkaProps.setProperty("group.id", config.get("kafka-groupid"));
long maxOutOfOrderness = 5 * 1000L;
SingleOutputStreamOperator<UserClickModel> dataStream = env.addSource(
new FlinkKafkaConsumer<>(
config.get("kafka-topic"),
new SimpleStringSchema(),
kafkaProps
))
.assignTimestampsAndWatermarks(WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofMillis(maxOutOfOrderness))
.withTimestampAssigner((element, recordTimestamp) -> {
return Long.valueOf(JSON.parseObject(element).getString("timestamp")) * 1000;
}))
.withIdleness(Duration.ofSeconds(1))
.map(new FCClickMapFunction()).returns(TypeInformation.of(new TypeHint<UserClickModel>() {}));
dataStream.keyBy(new KeySelector<UserClickModel, Tuple2<String, String>>() {
@Override
public Tuple2<String, String> getKey(UserClickModel value) throws Exception {
return Tuple2.of(value.getDate(), value.getProduct());
}
})
.window(TumblingEventTimeWindows.of(Time.days(1), Time.hours(-8)))
.trigger(ContinuousEventTimeTrigger.of(Time.seconds(10)))
.process(new MyProcessWindowFunctionBitMap())
.addSink(new FCClickSinkFunction());
env.execute(UserClickMain.class.getSimpleName());
}
}
java
import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import org.roaringbitmap.Roaring64NavigableMap;
public class MyProcessWindowFunctionBitMap extends ProcessWindowFunction<UserClickModel, UserClickModel, Tuple2<String, String>, TimeWindow> {
private transient ValueState<Integer> pvState;
private transient ValueState<Roaring64NavigableMap> bitMapState;
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
ValueStateDescriptor<Integer> pvStateDescriptor = new ValueStateDescriptor<>("pv", Integer.class);
ValueStateDescriptor<Roaring64NavigableMap> bitMapStateDescriptor = new ValueStateDescriptor("bitMap", TypeInformation.of(new TypeHint<Roaring64NavigableMap>() {}));
StateTtlConfig stateTtlConfig = StateTtlConfig
.newBuilder(Time.days(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
pvStateDescriptor.enableTimeToLive(stateTtlConfig);
bitMapStateDescriptor.enableTimeToLive(stateTtlConfig);
pvState = this.getRuntimeContext().getState(pvStateDescriptor);
bitMapState = this.getRuntimeContext().getState(bitMapStateDescriptor);
}
@Override
public void process(Tuple2<String, String> key, Context context, Iterable<UserClickModel> elements, Collector<UserClickModel> out) throws Exception {
Integer pv = pvState.value();
Roaring64NavigableMap bitMap = bitMapState.value();
if (bitMap == null) {
bitMap = new Roaring64NavigableMap();
pv = 0;
}
Iterator<UserClickModel> iterator = elements.iterator();
while (iterator.hasNext()) {
pv = pv + 1;
int uid = iterator.next().getUid();
bitMap.add(uid);
}
pvState.update(pv);
UserClickModel UserClickModel = new UserClickModel();
UserClickModel.setDate(key.f0);
UserClickModel.setProduct(key.f1);
UserClickModel.setPv(pv);
UserClickModel.setUv(bitMap.getIntCardinality());
out.collect(UserClickModel);
}
}
4: 要求每五分钟输出一次从凌晨到当前时间的统计值(类似GTV)
实现思路
从 keyBy 开始处理,设置1天的滑动窗口,步长为5分钟,在 process 中使用 if 判断数据是不是今天的来进行累加,这样过了00:00后,昨天的数据不会被统计,也就实现了业务要求的5分钟输出一次从凌晨到当前时间的统计值。
代码实现
java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
DataStreamSource<String> localSource = env.socketTextStream("localhost", 8888);
localSource.assignTimestampsAndWatermarks(
WatermarkStrategy.<Tuple3<String, Integer, String>>forBoundedOutOfOrderness(Duration.ZERO)
.withTimestampAssigner(new WyTimestampAssigner())
).keyBy(t -> t.getShop_name())
.timeWindow(Time.days(1), Time.minutes(5))
.process(new ProcessWindowFunction<GoodDetails, Tuple3<String, String, Integer>, String, TimeWindow>() {
SimpleDateFormat sdf_million = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS");
SimpleDateFormat sdf_day = new SimpleDateFormat("yyyy-MM-dd");
@Override
public void process(String s, Context ctx, Iterable<GoodDetails> elements, Collector<Tuple3<String, String, Integer>> out) throws Exception {
Calendar cld = Calendar.getInstance();
Iterator<GoodDetails> iterator = elements.iterator();
String curentDay = sdf_day.format(cld.getTimeInMillis() - 180000);
int countNum = 0;
while (iterator.hasNext()) {
GoodDetails next = iterator.next();
String elementData = next.getRegion_name().substring(0, 10);
if (elementData.equals(curentDay)) {
countNum += next.getGood_price();
}
}
long end = ctx.window().getEnd();
String windowEnd = sdf_million.format(end);
out.collect(Tuple3.of(windowEnd, s, countNum));
}
})
.name("sum-process").uid("sum-process");
env.execute();
5: 滑动窗口中,将数据分配到多个窗口
窗口的长度 / 窗口滑动的步长 = 窗口的个数
数据的流向和 TumblingEventTimeWindows 是一样的,所以直接跳到对应数据分配的地方 WindowOperator.processElement。
java
@Override
public void processElement(StreamRecord<IN> element) throws Exception {
final Collection<W> elementWindows = windowAssigner.assignWindows(element.getValue(), element.getTimestamp(), windowAssignerContext);
boolean isSkippedElement = true;
final K key = this.<K>getKeyedStateBackend().getCurrentKey();
if (windowAssigner instanceof MergingWindowAssigner) {
} else {
for (W window : elementWindows) {
if (isWindowLate(window)) {
continue;
}
isSkippedElement = false;
windowState.setCurrentNamespace(window);
windowState.add(element.getValue());
registerCleanupTimer(window);
}
}
}
自定义滑动窗口
java
import org.apache.flink.api.common.ExecutionConfig;
import org.apache.flink.api.common.typeutils.TypeSerializer;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.WindowAssigner;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.triggers.EventTimeTrigger;
import org.apache.flink.streaming.api.windowing.triggers.Trigger;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import java.util.ArrayList;
import java.util.Calendar;
import java.util.Collection;
import java.util.List;
public class MyEventTimeWindow extends WindowAssigner<Object, TimeWindow> {
private final long size;
private final long slide;
private final long offset;
protected MyEventTimeWindow(long size, long slide, long offset) {
this.size = size;
this.slide = slide;
this.offset = offset;
}
public static MyEventTimeWindow of(Time size, Time slide, Time offset) {
return new MyEventTimeWindow(size.toMilliseconds(), slide.toMilliseconds(), offset.toMilliseconds());
}
public static MyEventTimeWindow of(Time size, Time slide) {
return new MyEventTimeWindow(size.toMilliseconds(), slide.toMilliseconds(), 0L);
}
@Override
public Collection<TimeWindow> assignWindows(Object element, long timestamp, WindowAssignerContext windowAssignerContext) {
Calendar calendar = Calendar.getInstance();
calendar.setTimeInMillis(timestamp);
calendar.set(Calendar.HOUR_OF_DAY, 0);
calendar.set(Calendar.MINUTE, 0);
calendar.set(Calendar.SECOND, 0);
calendar.set(Calendar.MILLISECOND, 0);
long winStart = calendar.getTimeInMillis();
calendar.add(Calendar.DATE, 1);
long winEnd = calendar.getTimeInMillis() + 1;
String format = String.format("window的开始时间:%s,window的结束时间:%s", winStart, winEnd);
System.out.println(format);
long currentWindowEnd = TimeWindow.getWindowStartWithOffset(timestamp, this.offset, this.slide) + slide;
System.out.println(TimeWindow.getWindowStartWithOffset(timestamp, this.offset, this.slide) + "====" + currentWindowEnd);
int windowCounts = (int) ((winEnd - currentWindowEnd) / slide);
List<TimeWindow> windows = new ArrayList<>(windowCounts);
long currentEnd = currentWindowEnd;
if (timestamp > Long.MIN_VALUE) {
while (currentEnd < winEnd) {
windows.add(new TimeWindow(winStart, currentEnd));
currentEnd += slide;
}
}
return windows;
}
@Override
public Trigger<Object, TimeWindow> getDefaultTrigger(StreamExecutionEnvironment streamExecutionEnvironment) {
return EventTimeTrigger.create();
}
@Override
public TypeSerializer<TimeWindow> getWindowSerializer(ExecutionConfig executionConfig) {
return new TimeWindow.Serializer();
}
@Override
public boolean isEventTime() {
return true;
}
}
end