公司动态
Flink demo代码
Flink流任务用于checkpoint问题验证public class StreamDemo { public static void main(String[] args) throws Exception { // 1. 创建执行环境 // StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); // 2. 设置并行度 env.setParallelism(1); // 3. 启用Checkpoint使作业可以长期运行 env.enableCheckpointing(5000); // 5秒一次 env.getCheckpointConfig().setCheckpointStorage(file:////Users/wangqin/IdeaProjects/MyFlinkCode/checkpoint); CheckpointConfig checkpointConfig env.getCheckpointConfig(); checkpointConfig.setExternalizedCheckpointCleanup( CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION ); env.getCheckpointConfig().setTolerableCheckpointFailureNumber(0); // 4. 创建无限数据源一直产生数据 DataStreamString infiniteStream env.addSource(new InfiniteSource()).setParallelism(4); // 5. 简单的数据处理 infiniteStream .map(value - 处理数据: value , 时间: System.currentTimeMillis()) .print(); // 6. 执行作业会一直运行 env.execute(Infinite Stream Job); } /** * 自定义无限数据源 * 每秒生成一个随机数一直运行 */ public static class InfiniteSource extends RichParallelSourceFunctionString implements CheckpointedFunction { private volatile boolean isRunning true; private long count 0; Override public void run(SourceContextString ctx) throws Exception { while (isRunning) { // 每秒生成一个数据 String data 随机数- (int)(Math.random() * 100) -计数- (count); ctx.collect(data); // 控制数据生成速度每秒1条 Thread.sleep(1000); // 每10条数据打印一次日志 if (count % 10 0) { System.out.println([Source] 已生成 count 条数据继续运行中...); } } } Override public void cancel() { isRunning false; System.out.println(数据源已停止); } Override public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception { int subtaskIndex getRuntimeContext().getIndexOfThisSubtask(); System.out.println(subtaskIndex is subtaskIndex); // if (subtaskIndex 0) { // Thread.sleep(10000); // throw new RuntimeException(checkpoint failed); // } if (functionSnapshotContext.getCheckpointId() % 3 0) { throw new RuntimeException(checkpoint failed); } } Override public void initializeState(FunctionInitializationContext functionInitializationContext) throws Exception { } } }本地运行Flink任务发现在snapshotState方法中抛出异常不会生成checkpoint目录及metadata file延长方法的执行时间在执行checkpoint的时候会先生成空的checkpoint目录。checkpoint metadata file解析public class CheckpointMetaDataAnalyzer { public static void main(String[] args) throws Exception { File metadataFile new File(./checkpoint/d3299c5896dee8c5870606e41b0e0c0a/chk-5, _metadata); try (DataInputStream dis new DataInputStream(Files.newInputStream(metadataFile.toPath()))) { ClassLoader classLoader Thread.currentThread().getContextClassLoader(); CheckpointMetadata metadata Checkpoints.loadCheckpointMetadata( dis, classLoader, ./checkpoint/d3299c5896dee8c5870606e41b0e0c0a/chk-5); System.out.println(Loaded checkpoint ID: metadata.getCheckpointId()); } } }Flink watermark demopublic class WatermarkDemo { public static void main(String[] args) throws Exception { // 步骤1创建环境 StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); // 步骤2创建数据 DataStreamEvent stream env.addSource(new EventSource()); // 步骤3设置Watermark策略 WatermarkStrategyEvent strategy WatermarkStrategy .EventforBoundedOutOfOrderness(Duration.ofSeconds(1)) .withTimestampAssigner((event, recordTimestamp) - event.timestamp); // WatermarkStrategyEvent strategy WatermarkStrategy.noWatermarks(); // 步骤4应用Watermark策略 DataStreamEvent withWatermarks stream.assignTimestampsAndWatermarks(strategy); DataStreamEvent result withWatermarks.keyBy(event - event.id) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .reduce(new ReduceFunctionEvent() { Override public Event reduce(Event value1, Event value2) throws Exception { // 累加value值 return new Event(value1.id, Math.max(value1.timestamp, value2.timestamp), value1.value value2.value); } }); result.print(); env.execute(Watermark Example); } public static class Event { public String id; public long timestamp; public int value; public Event() {} public Event(String id, long timestamp, int value) { this.id id; this.timestamp timestamp; this.value value; } Override public String toString() { return String.format(Event{id%s, time%d, value%d}, id, timestamp, value); } } public static class EventSource implements SourceFunctionEvent { private volatile boolean isRunning true; private long count 0; Override public void run(SourceContextEvent ctx) throws Exception { while (isRunning) { Event eventA new Event(A, (count 3) * 1000L, 1); Event eventB new Event(B, (count 3) * 1000L, 2); Event eventC new Event(C, (count 3) * 1000L, 3); ctx.collectWithTimestamp(eventA, eventA.timestamp); ctx.collectWithTimestamp(eventB, eventB.timestamp); ctx.collectWithTimestamp(eventC, eventC.timestamp); Thread.sleep(5000); count 5; } } Override public void cancel() { isRunning false; System.out.println(数据源已停止); } } }Flink watermark demoflinksql模式public class WatermarkSqlDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String sourceDDL CREATE TABLE user_events (\n user_id BIGINT,\n event_time AS TIMESTAMPADD(SECOND, user_id, TO_TIMESTAMP(2026-01-01 00:00:00)),\n action STRING,\n WATERMARK FOR event_time AS event_time - INTERVAL 1 SECOND\n ) WITH (\n connector datagen,\n rows-per-second 1,\n fields.action.kind random,\n fields.action.length 5,\n fields.user_id.kind sequence,\n fields.user_id.start 0,\n fields.user_id.end 86400\n ); tableEnv.executeSql(sourceDDL); // 实时统计每5秒窗口的用户行为数 String windowQuery SELECT \n user_id,\n TUMBLE_START(event_time, INTERVAL 5 SECOND) as window_start,\n COUNT(*) as event_count,\n LISTAGG(action) as actions\n FROM user_events\n GROUP BY TUMBLE(event_time, INTERVAL 5 SECOND), user_id; // 创建打印输出的sink tableEnv.executeSql( CREATE TABLE window_results (\n user_id BIGINT,\n window_start TIMESTAMP(3),\n event_count BIGINT,\n actions STRING\n ) WITH (connector print) ); // 持续插入结果 tableEnv.executeSql(INSERT INTO window_results windowQuery); } }Flink处理时间demopublic class ProcessingTimeDemo { public static void main(String[] args) throws Exception { // 步骤1创建环境 StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); // 步骤2创建数据 DataStreamEvent stream env.addSource(new EventSource()); DataStreamEvent result stream.keyBy(event - event.id) .window(TumblingProcessingTimeWindows.of(Time.seconds(5))) .reduce(new ReduceFunctionEvent() { Override public Event reduce(Event value1, Event value2) throws Exception { // 累加value值 return new Event(value1.id, System.currentTimeMillis(), value1.value value2.value); } }); result.print(); env.execute(Watermark Example); } public static class Event { public String id; public long timestamp; public int value; public Event() {} public Event(String id, long timestamp, int value) { this.id id; this.timestamp timestamp; this.value value; } Override public String toString() { return String.format(Event{id%s, time%d, value%d}, id, timestamp, value); } } public static class EventSource implements SourceFunctionEvent { private volatile boolean isRunning true; private long count 0; Override public void run(SourceContextEvent ctx) throws Exception { while (isRunning) { Event eventA new Event(A, (count 3) * 1000L, 1); Event eventB new Event(B, (count 3) * 1000L, 2); Event eventC new Event(C, (count 3) * 1000L, 3); ctx.collect(eventA); ctx.collect(eventB); ctx.collect(eventC); Thread.sleep(5000); count 5; } } Override public void cancel() { isRunning false; System.out.println(数据源已停止); } } }Flink处理时间demoflinksql模式public class ProcessingTimeSqlDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String sourceDDL CREATE TABLE user_events (\n user_id BIGINT,\n action STRING,\n proc_time AS PROCTIME()\n ) WITH (\n connector datagen,\n rows-per-second 1,\n fields.action.kind random,\n fields.action.length 5,\n fields.user_id.kind sequence,\n fields.user_id.start 0,\n fields.user_id.end 86400\n ); tableEnv.executeSql(sourceDDL); // 使用处理时间进行窗口统计每5秒窗口的用户行为数 String windowQuery SELECT \n user_id,\n TUMBLE_START(proc_time, INTERVAL 5 SECOND) as window_start,\n COUNT(*) as event_count,\n LISTAGG(action) as actions\n FROM user_events\n GROUP BY TUMBLE(proc_time, INTERVAL 5 SECOND), user_id; // 创建打印输出的sink tableEnv.executeSql( CREATE TABLE window_results (\n user_id BIGINT,\n window_start TIMESTAMP(3),\n event_count BIGINT,\n actions STRING\n ) WITH (connector print) ); // 持续插入结果 tableEnv.executeSql(INSERT INTO window_results windowQuery); } }TVF窗口表值函数 demopublic class TVFDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String sourceDDL CREATE TABLE Bid (\n user_id BIGINT,\n price DECIMAL(10, 2),\n bid_time AS TIMESTAMPADD(SECOND, user_id, TO_TIMESTAMP(2026-01-01 00:00:00)),\n WATERMARK FOR bid_time AS bid_time - INTERVAL 1 SECOND\n ) WITH (\n connector datagen,\n rows-per-second 1,\n fields.price.min 10,\n fields.price.max 20,\n fields.user_id.kind sequence,\n fields.user_id.start 0,\n fields.user_id.end 86400\n ); tableEnv.executeSql(sourceDDL); // String windowQuery SELECT \n // window_start,\n // window_end,\n // AVG(price - 0.245) AS avg_price\n // FROM table (\n // TUMBLE(TABLE Bid, DESCRIPTOR(bid_time), INTERVAL 5 SECOND))\n // GROUP BY window_start, window_end; String windowQuery SELECT \n window_start,\n window_end,\n AVG(CASE \n WHEN price 0 \n THEN price - 0.245 \n END) AS avg_price\n FROM table (\n TUMBLE(TABLE Bid, DESCRIPTOR(bid_time), INTERVAL 5 SECOND))\n GROUP BY window_start, window_end; tableEnv.executeSql( CREATE TABLE window_results (\n window_start TIMESTAMP(3),\n window_end TIMESTAMP(3),\n avg_price DECIMAL(10, 2)\n ) WITH (connector print) ); tableEnv.executeSql(INSERT INTO window_results windowQuery); } }public class TVF2Demo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(2); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String sourceDDL CREATE TABLE Bid (\n user_id BIGINT,\n total DECIMAL(10, 2),\n num DECIMAL(10, 2),\n bid_time AS TIMESTAMPADD(SECOND, user_id, TO_TIMESTAMP(2026-01-01 00:00:00)),\n WATERMARK FOR bid_time AS bid_time - INTERVAL 1 SECOND\n ) WITH (\n connector datagen,\n rows-per-second 2,\n fields.total.min 10,\n fields.total.max 20,\n fields.num.min 1,\n fields.num.max 2,\n fields.user_id.kind sequence,\n fields.user_id.start 0,\n fields.user_id.end 86400\n ); tableEnv.executeSql(sourceDDL); // String windowQuery WITH BidData AS (\n // SELECT bid_time, total, num FROM Bid\n // )\n // SELECT \n // window_start,\n // window_end,\n // AVG(CASE \n // WHEN total 0 \n // AND num 0 \n // THEN total/num \n // END) AS avg_price,\n // AVG(CASE \n // WHEN total 0 \n // THEN total - 10 \n // END) AS total_2,\n // AVG(CASE \n // WHEN total 0 \n // THEN num\n // END) AS num_2\n // FROM table (\n // TUMBLE(TABLE BidData, DESCRIPTOR(bid_time), INTERVAL 5 SECOND))\n // GROUP BY window_start, window_end; String windowQuery WITH BidData AS (\n SELECT bid_time, total, num FROM Bid\n )\n SELECT \n window_start,\n window_end,\n (CAST(AVG( CASE \n WHEN total 0 \n AND num 0 \n THEN total/num \n END ) AS DOUBLE)) AS avg_price,\n (CAST(AVG( CASE \n WHEN total 0 \n THEN total - 10 \n END ) AS DOUBLE)) AS total_2,\n (CAST(AVG( CASE \n WHEN total 0 \n THEN num\n END ) AS DOUBLE)) AS num_2\n FROM table (\n TUMBLE(TABLE BidData, DESCRIPTOR(bid_time), INTERVAL 5 SECOND))\n GROUP BY window_start, window_end; tableEnv.executeSql( CREATE TABLE window_results (\n window_start TIMESTAMP(3),\n window_end TIMESTAMP(3),\n avg_price DOUBLE,\n total_2 DOUBLE,\n num_2 DOUBLE\n ) WITH (connector print) ); tableEnv.executeSql(INSERT INTO window_results windowQuery); } }flinksql kafka demopublic static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String createSourceTable CREATE TABLE kafka_source (\n id INT,\n name STRING,\n age INT,\n bid_time TIMESTAMP(3),\n WATERMARK FOR bid_time AS bid_time - INTERVAL 1 SECOND\n ) WITH (\n connector kafka,\n topic test-topic,\n properties.bootstrap.servers localhost:9092,\n properties.group.id flink-group,\n format csv,\n scan.startup.mode earliest-offset\n ); tableEnv.executeSql(createSourceTable); String windowQuery SELECT \n window_start,\n window_end,\n id,\n name,\n AVG(age) AS avg_age\n FROM table (\n TUMBLE(TABLE kafka_source, DESCRIPTOR(bid_time), INTERVAL 5 SECOND))\n GROUP BY window_start, window_end, id, name; tableEnv.executeSql( CREATE TABLE window_results (\n window_start TIMESTAMP(3),\n window_end TIMESTAMP(3),\n id INT,\n name STRING,\n avg_age DECIMAL(10, 2)\n ) WITH (connector print) ); tableEnv.executeSql(INSERT INTO window_results windowQuery); }flinksql kafka设置Idlenesspublic static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(2); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String createSourceTable CREATE TABLE kafka_source (\n id INT,\n name STRING,\n age INT,\n row_time BIGINT\n ) WITH (\n connector kafka,\n topic test-topic-2,\n properties.bootstrap.servers localhost:9092,\n properties.group.id flink-group,\n format csv,\n scan.startup.mode earliest-offset\n ); tableEnv.executeSql(createSourceTable); Table sourceTable tableEnv.from(kafka_source); DataStreamRow dataStream tableEnv.toDataStream(sourceTable); DataStreamProcessedRecord processedStream dataStream.map(row - { Long timestamp row.getFieldAs(row_time); return new ProcessedRecord( row.getFieldAs(id), row.getFieldAs(name), row.getFieldAs(age), timestamp ); }); WatermarkStrategyProcessedRecord strategy WatermarkStrategy .ProcessedRecordforBoundedOutOfOrderness(Duration.ofSeconds(1)) .withTimestampAssigner((event, recordTimestamp) - event.getTimestamp()) .withIdleness(Duration.ofSeconds(15)); DataStreamProcessedRecord dataStreamWithWatermarks processedStream.assignTimestampsAndWatermarks(strategy); ListExpression expressionList Stream.of(id, name, age) .map(Expressions::$) .collect(Collectors.toList()); expressionList.add(Expressions.$(row_time).rowtime()); tableEnv.createTemporaryView(processed_source, dataStreamWithWatermarks, expressionList.toArray(new Expression[0])); String windowQuery SELECT \n window_start,\n window_end,\n id,\n name,\n AVG(age) AS avg_age\n FROM table (\n TUMBLE(TABLE processed_source, DESCRIPTOR(row_time), INTERVAL 5 SECOND))\n GROUP BY window_start, window_end, id, name; tableEnv.executeSql( CREATE TABLE window_results (\n window_start TIMESTAMP(3),\n window_end TIMESTAMP(3),\n id INT,\n name STRING,\n avg_age DECIMAL(10, 2)\n ) WITH (connector print) ); tableEnv.executeSql(INSERT INTO window_results windowQuery); } public static class ProcessedRecord { public int id; public String name; public int age; public long timestamp; public ProcessedRecord() {} public ProcessedRecord(int id, String name, int age, long timestamp) { this.id id; this.name name; this.age age; this.timestamp timestamp; } public long getTimestamp() { return timestamp; } }flinksql mysql demopublic static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .createLocalEnvironmentWithWebUI(new Configuration()); env.setParallelism(1); env.disableOperatorChaining(); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); String jdbcUrl jdbc:mysql://localhost:3306/mysql_test ?useSSLfalse allowPublicKeyRetrievaltrue serverTimezoneAsia/Shanghai; tableEnv.executeSql( CREATE TABLE source_users ( user_id INT, username STRING, age INT ) WITH ( connector jdbc, url jdbcUrl , table-name users_in, username root, password root ) ); tableEnv.executeSql( CREATE TABLE target_users ( user_id INT, username STRING, age INT ) WITH ( connector jdbc, url jdbcUrl , table-name users_out, username root, password root ) ); tableEnv.executeSql( INSERT INTO target_users SELECT * FROM source_users ); }