时间:2026-03-01 15:03
人气:
作者:admin
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class EventTimeStreamDemo { public static void main(String[] args) throws Exception { // 1. 创建执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 开启事件时间(Flink 1.12+ 默认开启,但显式声明更规范) env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 2. 配置Kafka Source,读取订单流(事件时间流) KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") // Kafka集群地址 .setTopics("order_topic") // 订阅的订单主题 .setGroupId("flink_order_group") // 消费者组 // 从最新偏移量开始读取(生产环境可根据需求调整为 earliest) .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) // 字符串反序列化 .build(); // 3. 读取Kafka数据,指定事件时间字段(假设订单数据格式:orderId,eventTime,amount) DataStream<Order> orderStream = env.fromSource( kafkaSource, // 水位线策略:基于事件时间字段,允许3秒乱序(后续水位线章节详细说明) WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(3)) .mapTimestamp(line -> { // 解析订单数据,提取事件时间戳(毫秒级) String[] fields = line.split(","); return Long.parseLong(fields[1]); }), "Kafka Order Source" ) // 将字符串转换为Order实体类 .map(line -> { String[] fields = line.split(","); return new Order( fields[0], Long.parseLong(fields[1]), Double.parseDouble(fields[2]) ); }); // 后续可对orderStream进行窗口、聚合等操作 orderStream.print("Event Time Order Stream"); // 执行任务 env.execute("Flink Event Time Stream Demo"); } // 订单实体类 static class Order { private String orderId; private Long eventTime; // 事件时间戳(毫秒) private Double amount; // 构造方法、getter/setter省略 public Order(String orderId, Long eventTime, Double amount) { this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } @Override public String toString() { return "Order{orderId='" + orderId + "', eventTime=" + eventTime + ", amount=" + amount + "}"; } } }
说明:该示例创建了基于Kafka的事件时间流,核心是通过
WatermarkStrategy指定事件时间字段,并设置3秒乱序容忍,为后续水位线和窗口计算奠定基础;同时使用Operator State(Kafka偏移量状态),Flink会自动维护偏移量,避免重复读取。
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.functions.AggregateFunction; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; public class WindowDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 1. 读取订单流(复用上面的orderStream,此处简化) DataStream<Order> orderStream = getOrderStream(env); // 2. 滚动窗口:每10分钟统计一次订单总数和总金额(事件时间) DataStream<OrderStats> tumblingWindowResult = orderStream // 按窗口ID分组(此处无需额外分组,窗口本身按时间划分) .windowAll(TumblingEventTimeWindows.of(Time.minutes(10))) // 聚合计算:统计订单数和总金额 .aggregate(new OrderAggregateFunction()); // 3. 滑动窗口:每5分钟统计一次过去10分钟的订单数据(事件时间) DataStream<OrderStats> slidingWindowResult = orderStream .windowAll(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(5))) .aggregate(new OrderAggregateFunction()); // 输出结果 tumblingWindowResult.print("滚动窗口(10分钟)统计结果"); slidingWindowResult.print("滑动窗口(10分钟窗口,5分钟滑动)统计结果"); env.execute("Flink Window Demo"); } // 聚合函数:统计每个窗口的订单总数和总金额 static class OrderAggregateFunction implements AggregateFunction<Order, OrderStats, OrderStats> { // 初始化聚合状态(初始订单数0,总金额0) @Override public OrderStats createAccumulator() { return new OrderStats(0L, 0.0); } // 累加数据:每来一条订单,更新状态 @Override public OrderStats add(Order order, OrderStats accumulator) { return new OrderStats( accumulator.getOrderCount() + 1, accumulator.getTotalAmount() + order.getAmount() ); } // 窗口触发时,输出聚合结果 @Override public OrderStats getResult(OrderStats accumulator) { return accumulator; } // 并行窗口的状态合并(windowAll无需合并,多并行时需实现) @Override public OrderStats merge(OrderStats a, OrderStats b) { return new OrderStats( a.getOrderCount() + b.getOrderCount(), a.getTotalAmount() + b.getTotalAmount() ); } } // 订单统计结果实体类 static class OrderStats { private Long orderCount; // 订单总数 private Double totalAmount; // 订单总金额 // 构造方法、getter/setter省略 public OrderStats(Long orderCount, Double totalAmount) { this.orderCount = orderCount; this.totalAmount = totalAmount; } @Override public String toString() { return "OrderStats{orderCount=" + orderCount + ", totalAmount=" + totalAmount + "}"; } // getter方法 public Long getOrderCount() { return orderCount; } public Double getTotalAmount() { return totalAmount; } } // 简化:获取订单流(实际可复用代码示例1的Kafka Source逻辑) private static DataStream<Order> getOrderStream(StreamExecutionEnvironment env) { // 模拟订单数据(实际替换为Kafka Source) return env.fromElements( new Order("1001", 1683000625000L, 99.0), // 2024-05-01 10:03:45 new Order("1002", 1683001225000L, 199.0),// 2024-05-01 10:10:25 new Order("1003", 1683001825000L, 299.0) // 2024-05-01 10:20:25 ) // 模拟水位线生成(后续章节详细说明) .assignTimestampsAndWatermarks( WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofMinutes(5)) .withTimestampAssigner((order, timestamp) -> order.getEventTime()) ); } // 复用Order实体类(同代码示例1) static class Order { private String orderId; private Long eventTime; private Double amount; public Order(String orderId, Long eventTime, Double amount) { this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } public Long getEventTime() { return eventTime; } public Double getAmount() { return amount; } } }
说明:该示例实现了滚动窗口和滑动窗口的核心逻辑,通过
AggregateFunction实现订单数和总金额的聚合,窗口的触发由后续的水位线控制;聚合过程中,中间结果会自动存储在Window State(Keyed State的一种)中,无需手动管理。
import org.apache.flink.api.common.eventtime.*; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.util.OutputTag; public class WatermarkDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); // 并行度设置为1(方便测试,生产环境根据集群配置调整) env.setParallelism(1); // 1. 定义迟到数据输出标签(用于收集窗口关闭后到达的迟到数据) OutputTag<Order> lateDataTag = new OutputTag<Order>("late_order_data"){}; // 2. 读取订单流,生成水位线 DataStream<Order> orderStream = env.fromElements( new Order("1001", 1683000000000L, 99.0), // 10:00:00 new Order("1002", 1683000599000L, 199.0),// 10:09:59(窗口内最后一条正常数据) new Order("1003", 1683000601000L, 299.0),// 10:10:01(迟到1秒) new Order("1004", 1683000900000L, 399.0) // 10:15:00(迟到5分钟,超过允许迟到时间) ) // 生成水位线:允许5分钟乱序(对应场景中的允许迟到时间) .assignTimestampsAndWatermarks( new WatermarkStrategy<Order>() { @Override public WatermarkGenerator<Order> createWatermarkGenerator(WatermarkGeneratorSupplier.Context context) { // 周期性水位线生成器:每100ms生成一次水位线 return new PeriodicWatermarkGenerator<Order>() { // 当前最大事件时间 private long maxEventTime = Long.MIN_VALUE; // 允许迟到时间(5分钟,转换为毫秒) private final long allowedLateness = 5 * 60 * 1000; @Override public void onEvent(Order order, long eventTimestamp, WatermarkOutput output) { // 每接收一条事件,更新最大事件时间 maxEventTime = Math.max(maxEventTime, eventTimestamp); } @Override public void onPeriodicEmit(WatermarkOutput output) { // 生成水位线:当前最大事件时间 - 允许迟到时间 Watermark watermark = new Watermark(maxEventTime - allowedLateness); output.emitWatermark(watermark); } }; } } // 指定事件时间字段(Order类的eventTime属性) .withTimestampAssigner((order, timestamp) -> order.getEventTime()) ); // 3. 滚动窗口(10分钟),处理迟到数据 SingleOutputStreamOperator<OrderStats> windowResult = orderStream .windowAll(TumblingEventTimeWindows.of(Time.minutes(10))) // 设置允许迟到时间(5分钟),与水位线策略一致 .allowedLateness(Time.minutes(5)) // 将超过允许迟到时间的迟到数据,输出到侧输出流 .sideOutputLateData(lateDataTag) // 聚合计算 .aggregate(new OrderAggregateFunction()); // 4. 输出窗口计算结果和迟到数据 windowResult.print("窗口计算结果"); // 读取侧输出流的迟到数据(可用于后续补算) windowResult.getSideOutput(lateDataTag).print("迟到数据(超过5分钟)"); env.execute("Flink Watermark & Late Data Demo"); } // 复用聚合函数和实体类(同代码示例2) static class OrderAggregateFunction implements AggregateFunction<Order, OrderStats, OrderStats> { @Override public OrderStats createAccumulator() { return new OrderStats(0L, 0.0); } @Override public OrderStats add(Order order, OrderStats accumulator) { return new OrderStats(accumulator.getOrderCount() + 1, accumulator.getTotalAmount() + order.getAmount()); } @Override public OrderStats getResult(OrderStats accumulator) { return accumulator; } @Override public OrderStats merge(OrderStats a, OrderStats b) { return new OrderStats(a.getOrderCount() + b.getOrderCount(), a.getTotalAmount() + b.getTotalAmount()); } } static class OrderStats { private Long orderCount; private Double totalAmount; public OrderStats(Long orderCount, Double totalAmount) { this.orderCount = orderCount; this.totalAmount = totalAmount; } @Override public String toString() { return "OrderStats{orderCount=" + orderCount + ", totalAmount=" + totalAmount + "}"; } public Long getOrderCount() { return orderCount; } public Double getTotalAmount() { return totalAmount; } } static class Order { private String orderId; private Long eventTime; private Double amount; public Order(String orderId, Long eventTime, Double amount) { this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } public Long getEventTime() { return eventTime; } public Double getAmount() { return amount; } @Override public String toString() { return "Order{orderId='" + orderId + "', eventTime=" + eventTime + ", amount=" + amount + "}"; } } }
说明:该示例实现了自定义水位线生成器,明确了“水位线=当前最大事件时间-允许迟到时间”的核心逻辑;同时通过
allowedLateness设置窗口允许迟到时间,通过侧输出流收集超过允许迟到时间的数据,解决了“数据迟到”的核心痛点。
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.state.ValueState; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.api.common.time.Time; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.KeyedProcessFunction; import org.apache.flink.util.Collector; public class KeyedStateDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.setParallelism(1); // 1. 读取订单流(按用户ID分组,统计每个用户的订单总额) DataStream<Order> orderStream = env.fromElements( new Order("user1", "1001", 1683000625000L, 99.0), new Order("user1", "1002", 1683001225000L, 199.0), new Order("user2", "1003", 1683001825000L, 299.0), new Order("user1", "1004", 1683002425000L, 399.0) ) .assignTimestampsAndWatermarks( WatermarkStrategy.<Order>forBoundedOutOfOrderness(Time.seconds(3)) .withTimestampAssigner((order, timestamp) -> order.getEventTime()) ); // 2. 按用户ID分组,使用Keyed State统计每个用户的订单总额 DataStream<UserOrderTotal> userTotalStream = orderStream .keyBy(Order::getUserId) // 按用户ID分组,每个用户对应一个独立的状态实例 .process(new KeyedProcessFunction<String, Order, UserOrderTotal>() { // 定义Keyed State:存储当前用户的订单总额(ValueState是最常用的Keyed State类型) private ValueState<Double> userTotalAmountState; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 初始化状态,设置状态TTL(过期时间):1小时未更新则自动清理 ValueStateDescriptor<Double> stateDescriptor = new ValueStateDescriptor<>( "user_total_amount", // 状态名称 Double.class // 状态类型 ); // 配置状态TTL StateTtlConfig ttlConfig = StateTtlConfig.newBuilder(Time.hours(1)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) // 创建/更新时刷新TTL .setStateVisibility(StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp) .build(); stateDescriptor.enableTimeToLive(ttlConfig); // 获取状态实例 userTotalAmountState = getRuntimeContext().getState(stateDescriptor); } @Override public void processElement(Order order, Context ctx, Collector<UserOrderTotal> out) throws Exception { // 读取当前状态中的订单总额(若状态未初始化,默认值为null) Double currentTotal = userTotalAmountState.value(); if (currentTotal == null) { currentTotal = 0.0; } // 更新状态:累加当前订单金额 currentTotal += order.getAmount(); userTotalAmountState.update(currentTotal); // 输出当前用户的订单总额 out.collect(new UserOrderTotal(order.getUserId(), currentTotal)); } }); // 输出结果 userTotalStream.print("每个用户订单总额统计"); env.execute("Flink Keyed State Demo"); } // 订单实体类(新增userId字段) static class Order { private String userId; private String orderId; private Long eventTime; private Double amount; public Order(String userId, String orderId, Long eventTime, Double amount) { this.userId = userId; this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } public String getUserId() { return userId; } public Long getEventTime() { return eventTime; } public Double getAmount() { return amount; } } // 用户订单总额实体类 static class UserOrderTotal { private String userId; private Double totalAmount; public UserOrderTotal(String userId, Double totalAmount) { this.userId = userId; this.totalAmount = totalAmount; } @Override public String toString() { return "UserOrderTotal{userId='" + userId + "', totalAmount=" + totalAmount + "}"; } } }
说明:该示例使用Keyed State(ValueState)统计每个用户的订单总额,核心是通过
ValueStateDescriptor初始化状态,并配置状态TTL(1小时),避免过期状态占用资源;每个用户ID对应一个独立的状态实例,实现了“按Key独立统计”的需求。
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.restartstrategy.RestartStrategies; import org.apache.flink.api.common.time.Time; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.runtime.state.filesystem.FsStateBackend; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; import java.util.concurrent.TimeUnit; public class CheckpointDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.setParallelism(1); // 1. 配置Checkpoint(核心:持久化状态,保障故障恢复) // 1.1 开启Checkpoint,间隔1分钟(1000ms * 60) env.enableCheckpointing(60000); // 1.2 配置Checkpoint存储介质:HDFS(生产环境推荐),本地测试可用file:///tmp/flink-checkpoint env.setStateBackend(new FsStateBackend("hdfs://localhost:9000/flink/checkpoints")); // 1.3 配置Checkpoint参数 env.getCheckpointConfig().setCheckpointTimeout(30000); // Checkpoint超时时间:30秒 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); // 两次Checkpoint最小间隔:30秒 env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3); // 允许Checkpoint失败次数:3次 // 1.4 配置故障重启策略:失败后自动重启,最多重启3次,每次间隔5秒 env.setRestartStrategy(RestartStrategies.fixedDelayRestart( 3, // 最大重启次数 Time.of(5, TimeUnit.SECONDS) // 重启间隔 )); // 2. 配置Kafka Source(Operator State:偏移量由Checkpoint管理) KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setTopics("order_topic") .setGroupId("flink_order_checkpoint_group") // 从Checkpoint中恢复偏移量(若没有Checkpoint,从最新偏移量开始) .setStartingOffsets(OffsetsInitializer.restoreFromCheckpoint()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); // 3. 读取订单流,生成水位线 DataStream<Order> orderStream = env.fromSource( kafkaSource, WatermarkStrategy.<String>forBoundedOutOfOrderness(Time.minutes(5)) .mapTimestamp(line -> { String[] fields = line.split(","); return Long.parseLong(fields[2]); // 假设第三列为事件时间戳 }), "Kafka Source With Checkpoint" ) .map(line -> { String[] fields = line.split(","); return new Order(fields[0], fields[1], Long.parseLong(fields[2]), Double.parseDouble(fields[3])); }); // 4. 滚动窗口计算,状态由Checkpoint持久化 DataStream<OrderStats> windowResult = orderStream .windowAll(TumblingEventTimeWindows.of(Time.minutes(10))) .allowedLateness(Time.minutes(5)) .aggregate(new OrderAggregateFunction()); // 输出结果 windowResult.print("Checkpoint Demo 窗口计算结果"); env.execute("Flink Checkpoint & Fault Recovery Demo"); } // 复用聚合函数和实体类(同前面示例) static class OrderAggregateFunction implements AggregateFunction<Order, OrderStats, OrderStats> { @Override public OrderStats createAccumulator() { return new OrderStats(0L, 0.0); } @Override public OrderStats add(Order order, OrderStats accumulator) { return new OrderStats(accumulator.getOrderCount() + 1, accumulator.getTotalAmount() + order.getAmount()); } @Override public OrderStats getResult(OrderStats accumulator) { return accumulator; } @Override public OrderStats merge(OrderStats a, OrderStats b) { return new OrderStats(a.getOrderCount() + b.getOrderCount(), a.getTotalAmount() + b.getTotalAmount()); } } static class OrderStats { private Long orderCount; private Double totalAmount; public OrderStats(Long orderCount, Double totalAmount) { this.orderCount = orderCount; this.totalAmount = totalAmount; } @Override public String toString() { return "OrderStats{orderCount=" + orderCount + ", totalAmount=" + totalAmount + "}"; } public Long getOrderCount() { return orderCount; } public Double getTotalAmount() { return totalAmount; } } static class Order { private String userId; private String orderId; private Long eventTime; private Double amount; public Order(String userId, String orderId, Long eventTime, Double amount) { this.userId = userId; this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } public Long getEventTime() { return eventTime; } public Double getAmount() { return amount; } } }
说明:该示例完整配置了Checkpoint,包括存储介质(HDFS)、触发间隔、超时时间、重启策略等核心参数;Kafka Source通过
OffsetsInitializer.restoreFromCheckpoint()从Checkpoint中恢复偏移量,窗口聚合状态也会被定期持久化。当任务故障重启时,会从最近的Checkpoint快照中恢复所有状态(偏移量、聚合结果),实现“exactly-once”语义。
import org.apache.flink.api.common.eventtime.*; import org.apache.flink.api.common.functions.AggregateFunction; import org.apache.flink.api.common.restartstrategy.RestartStrategies; import org.apache.flink.api.common.state.StateTtlConfig; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.api.common.time.Time; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.runtime.state.filesystem.FsStateBackend; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows; import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.util.OutputTag; import java.util.concurrent.TimeUnit; /** * 五大组件完整协作示例:流(Kafka)+水位线+窗口+状态+Checkpoint * 功能:实时统计每10分钟的订单总数和总金额,允许5分钟迟到,支持故障恢复 */ public class FlinkFullCooperationDemo { public static void main(String[] args) throws Exception { // 1. 初始化执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.setParallelism(2); // 模拟分布式环境,多并行度 // 2. 配置Checkpoint(保障状态可靠) env.enableCheckpointing(60000); // 每1分钟触发一次Checkpoint env.setStateBackend(new FsStateBackend("hdfs://localhost:9000/flink/full-cooperation-checkpoints")); env.getCheckpointConfig().setCheckpointTimeout(30000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.of(5, TimeUnit.SECONDS))); // 3. 配置Kafka Source(流:数据入口) KafkaSource<String> kafkaSource = KafkaSource.<String>builder() .setBootstrapServers("localhost:9092") .setTopics("order_topic") .setGroupId("flink_full_cooperation_group") .setStartingOffsets(OffsetsInitializer.restoreFromCheckpoint()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); // 4. 读取流数据,生成水位线(时间标尺) OutputTag<Order> lateDataTag = new OutputTag<Order>("late_order"){}; DataStream<Order> orderStream = env.fromSource( kafkaSource, // 水位线策略:允许5分钟乱序 WatermarkStrategy.<String>forBoundedOutOfOrderness(Time.minutes(5)) .mapTimestamp(line -> { String[] fields = line.split(","); return Long.parseLong(fields[2]); // 第三列为事件时间戳(毫秒) }), "Kafka Order Source" ) .map(line -> { String[] fields = line.split(","); return new Order( fields[0], // userId fields[1], // orderId Long.parseLong(fields[2]), // eventTime Double.parseDouble(fields[3]) // amount ); }); // 5. 窗口(切割数据)+ 状态(存储中间结果)+ 聚合计算 SingleOutputStreamOperator<OrderWindowStats> windowResult = orderStream // 按窗口ID分组(此处用windowAll,多并行可用keyBy+window) .windowAll(TumblingEventTimeWindows.of(Time.minutes(10))) .allowedLateness(Time.minutes(5)) // 允许5分钟迟到 .sideOutputLateData(lateDataTag) // 收集超期迟到数据 .aggregate(new OrderWindowAggregate()); // 6. 输出结果 windowResult.print("每10分钟订单统计结果"); windowResult.getSideOutput(lateDataTag).print("超期迟到订单(补算用)"); // 7. 执行任务 env.execute("Flink 五大组件完整协作示例"); } // 窗口聚合函数:状态自动存储中间结果(Window State) static class OrderWindowAggregate implements AggregateFunction<Order, OrderWindowStats, OrderWindowStats> { // 初始化聚合状态(订单数0,总金额0) @Override public OrderWindowStats createAccumulator() { return new OrderWindowStats(0L, 0.0); } // 累加数据,更新状态 @Override public OrderWindowStats add(Order order, OrderWindowStats accumulator) { return new OrderWindowStats( accumulator.getOrderCount() + 1, accumulator.getTotalAmount() + order.getAmount() ); } // 窗口触发(水位线到达),输出结果 @Override public OrderWindowStats getResult(OrderWindowStats accumulator) { return accumulator; } // 多并行窗口状态合并 @Override public OrderWindowStats merge(OrderWindowStats a, OrderWindowStats b) { return new OrderWindowStats( a.getOrderCount() + b.getOrderCount(), a.getTotalAmount() + b.getTotalAmount() ); } } // 窗口统计结果实体类 static class OrderWindowStats { private Long orderCount; private Double totalAmount; public OrderWindowStats(Long orderCount, Double totalAmount) { this.orderCount = orderCount; this.totalAmount = totalAmount; } @Override public String toString() { return "OrderWindowStats{orderCount=" + orderCount + ", totalAmount=" + totalAmount + "}"; } public Long getOrderCount() { return orderCount; } public Double getTotalAmount() { return totalAmount; } } // 订单实体类 static class Order { private String userId; private String orderId; private Long eventTime; private Double amount; public Order(String userId, String orderId, Long eventTime, Double amount) { this.userId = userId; this.orderId = orderId; this.eventTime = eventTime; this.amount = amount; } public Long getEventTime() { return eventTime; } public Double getAmount() { return amount; } } }
说明:该示例是五大组件的完整协作实现,涵盖了“Kafka流(数据入口)→水位线(时间标尺)→滚动窗口(切割数据)→Window State(存储中间结果)→Checkpoint(持久化状态)”的全流程,与前文“实时统计每10分钟订单量”的场景完全对应,可直接用于生产环境参考;同时包含迟到数据处理、故障重启策略,贴合实际业务需求。
【从0到1构建一个ClaudeAgent】协作-Agent团队