package org.bigdatatechcir.learn_flink.part5_flink_watermark;import org.apache.flink.api.common.eventtime.WatermarkStrategy;import org.apache.flink.api.common.functions.MapFunction;import org.apache.flink.api.common.typeinfo.Types;import org.apache.flink.api.java.tuple.Tuple2;import org.apache.flink.api.java.tuple.Tuple3;import org.apache.flink.configuration.Configuration;import org.apache.flink.configuration.RestOptions;import org.apache.flink.streaming.api.TimeCharacteristic;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;import org.apache.flink.streaming.api.watermark.Watermark;import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;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.time.Duration;import java.time.Instant;import java.time.ZoneId;import java.time.ZonedDateTime;import java.time.format.DateTimeFormatter;import java.util.Random;public class WithIdLenessDemo {public static void main(String[] args) throws Exception {Configuration conf = new Configuration();conf.setString(RestOptions.BIND_PORT, "8081");final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(conf);env.setParallelism(1);env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);DataStream<String> text = env.addSource(new RichParallelSourceFunction<String>() {private volatile boolean running = true;private volatile long count = 0; // 计数器用于跟踪已生成的数据条数private final Random random = new Random();@Overridepublic void run(SourceContext<String> ctx) throws Exception {while (running) {int randomNum = random.nextInt(5) + 1;long timestamp = System.currentTimeMillis();ctx.collectWithTimestamp("key" + randomNum + "," + 1 + "," + timestamp, timestamp);if (++count % 200 == 0) { // 每200条数据发送一次Watermarkctx.emitWatermark(new Watermark(timestamp));System.out.println("Manual Watermark emitted: " + timestamp);}ZonedDateTime generateDataDateTime = Instant.ofEpochMilli(timestamp).atZone(ZoneId.systemDefault());DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS");String formattedGenerateDataDateTime = generateDataDateTime.format(formatter);System.out.println("Generated data: " + "key" + randomNum + "," + 1 + "," + timestamp + " at " + formattedGenerateDataDateTime);Thread.sleep(1000);}}@Overridepublic void cancel() {running = false;}});DataStream<Tuple3<String, Integer, Long>> tuplesWithTimestamp = text.map(new MapFunction<String, Tuple3<String, Integer, Long>>() {@Overridepublic Tuple3<String, Integer, Long> map(String value) {String[] words = value.split(",");return new Tuple3<>(words[0], Integer.parseInt(words[1]), Long.parseLong(words[2]));}}).returns(Types.TUPLE(Types.STRING, Types.INT, Types.LONG));// 设置 Watermark 策略DataStream<Tuple3<String, Integer, Long>> withWatermarks = tuplesWithTimestamp.assignTimestampsAndWatermarks(WatermarkStrategy.<Tuple3<String, Integer, Long>>forBoundedOutOfOrderness(Duration.ofSeconds(5))//处理空闲数据源.withIdleness(Duration.ofSeconds(15)).withTimestampAssigner((element, recordTimestamp) -> element.f2));// 窗口逻辑DataStream<Tuple2<String, Integer>> keyedStream = withWatermarks.keyBy(value -> value.f0).window(TumblingEventTimeWindows.of(Time.seconds(5))).process(new ProcessWindowFunction<Tuple3<String, Integer, Long>, Tuple2<String, Integer>, String, TimeWindow>() {@Overridepublic void process(String s, Context context, Iterable<Tuple3<String, Integer, Long>> elements, Collector<Tuple2<String, Integer>> out) throws Exception {int count = 0;for (Tuple3<String, Integer, Long> element : elements) {count++;}long start = context.window().getStart();long end = context.window().getEnd();ZonedDateTime startDateTime = Instant.ofEpochMilli(start).atZone(ZoneId.systemDefault());ZonedDateTime endDateTime = Instant.ofEpochMilli(end).atZone(ZoneId.systemDefault());DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS");String formattedStart = startDateTime.format(formatter);String formattedEnd = endDateTime.format(formatter);System.out.println("Tumbling Window [start " + formattedStart + ", end " + formattedEnd + ") for key " + s);// 输出窗口结束时的Watermarklong windowEndWatermark = context.currentWatermark();ZonedDateTime windowEndDateTime = Instant.ofEpochMilli(windowEndWatermark).atZone(ZoneId.systemDefault());String formattedWindowEndWatermark = windowEndDateTime.format(formatter);System.out.println("Watermark at the end of window: " + formattedWindowEndWatermark);out.collect(new Tuple2<>(s, count));}});// 输出结果keyedStream.print();// 执行任务env.execute("With Id Leness Demo");}}
这或许是一个对你有用的开源项目,data-warehouse-learning 项目是一套基于 MySQL + Kafka + Hadoop + Hive + Dolphinscheduler + Doris + Seatunnel + Paimon + Hudi + Iceberg + Flink + Dinky + DataRT + SuperSet 实现的实时离线数仓(数据湖)系统,以大家最熟悉的电商业务为切入点,详细讲述并实现了数据产生、同步、数据建模、数仓(数据湖)建设、数据服务、BI报表展示等数据全链路处理流程。
https://gitee.com/wzylzjtn/data-warehouse-learning
https://github.com/Mrkuhuo/data-warehouse-learning
https://bigdatacircle.top/
项目演示:
03
代码获取
https://gitee.com/wzylzjtn/data-warehouse-learning
https://github.com/Mrkuhuo/data-warehouse-learning

04
文档获取


05
进交流群群添加作者

推荐阅读
【视频】| 2024最新版【实时数仓数据湖实战教程视频】重磅更新,打通新手到大牛的最后一公里 【视频】| 1. Flink基本架构概述 【视频】| 2. Flink任务运行时架构图详解 【视频】| 3. Flink 任务、子任务、算子链及Slot共享详解 【视频】| 4. Flink四种API详解及代码演示 【视频】| 5. 如何在IDEA中启动Flink任务并看到WEBUI呢? 【视频】| 6. Centos8 安装Flink1.18.1 详细教程 【视频】| 7. Flink 使用 Session & Application 模式提交任务实战演示 【视频】| 8. Flink 三种时间语义详解及代码实战 【视频】| 9. Flink 窗口函数是个好东西,你真的会用吗? 【视频】| 10. Flink乱序问题解决神器Watermark详解及代码实战 【视频】| 11. Flink Watermark何种情况下可以触发窗口计算? 【视频】| 12.Flink 并行运行时Watermark如何向下传递? 【视频】| 13.Flink Watermark的两种产生方式:标点水位线和周期水位线详解
文章转载自大数据技能圈,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。





