IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 大数据 -> Flink idea测试checkPoint -> 正文阅读

[大数据]Flink idea测试checkPoint

在Flink StreamExecutionEnvironment对应的configuration中新增配置execution.savepoint.path就可以在启动Flink任务的时候默认从上一次的状态中恢复过来

        Configuration configuration1 = new Configuration();
        //flink parallelism=16 savepoint state
//        configuration1.setString("execution.savepoint.path",
//                "file:///Users/wenbao/checkPoint/d8ca368f349922c36c20498e1bedb9e7/chk-1");
        //flink parallelism=3  savepoint state
        configuration1.setString("execution.savepoint.path",
                "file:///Users/wenbao/checkPoint/efc27bd1f33fcbd5eae9ab8431d64251/chk-1");
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(
                configuration1);

下面是测试hybird source从flink checkPoint中恢复的代码:
?

    public static void main(String[] args) throws Exception {
        
        Configuration configuration1 = new Configuration();
        //flink parallelism=16 savepoint state
//        configuration1.setString("execution.savepoint.path",
//                "file:///Users/wenbao/checkPoint/d8ca368f349922c36c20498e1bedb9e7/chk-1");
        //flink parallelism=3  savepoint state
        configuration1.setString("execution.savepoint.path",
                "file:///Users/wenbao/checkPoint/efc27bd1f33fcbd5eae9ab8431d64251/chk-1");
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(
                configuration1);
        env.setParallelism(4);
        env.enableCheckpointing(30000, CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setCheckpointTimeout(2 * 60 * 1000);
        env.getCheckpointConfig().setTolerableCheckpointFailureNumber(Integer.MAX_VALUE);
        env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
        env.getCheckpointConfig().enableExternalizedCheckpoints(
                ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
        env.setStateBackend(new FsStateBackend("file:///Users/wenbao/checkPoint"));
        
        Path path = Path.fromLocalFile(new File("a"));
        
        FileSource<String> source = FileSource.forRecordStreamFormat(new TextLineFormat(), path)
                .build();
        HybridSource<String> hybridSource =
                HybridSource.builder(source)
                        .addSource(
                                KafkaSource.<String>builder()
                                        .setStartingOffsets(
                                                OffsetsInitializer.timestamp(1648542406228L))
                                        .setBootstrapServers(
                                                "localhost:9092")
                                        .setTopics("wb_test")
                                        .setGroupId("wb-test")
                                        .setDeserializer(
                                                new KafkaRecordDeserializationSchema<String>() {
                                                    @Override
                                                    public void deserialize(
                                                            ConsumerRecord<byte[], byte[]> record,
                                                            Collector<String> out)
                                                            throws IOException {
                                                        out.collect(new String(record.value()));
                                                    }
                                                    
                                                    @Override
                                                    public TypeInformation<String> getProducedType() {
                                                        return TypeInformation.of(String.class);
                                                    }
                                                })
                                        .build())
                        .build();
        
        DataStreamSource<String> stringDataStreamSource = env.fromSource(hybridSource,
                WatermarkStrategy.noWatermarks(), "file-source", TypeInformation.of(String.class));
        
        stringDataStreamSource.print();
        
        env.execute("aaa");
        
    }

  大数据 最新文章
实现Kafka至少消费一次
亚马逊云科技:还在苦于ETL?Zero ETL的时代
初探MapReduce
【SpringBoot框架篇】32.基于注解+redis实现
Elasticsearch:如何减少 Elasticsearch 集
Go redis操作
Redis面试题
专题五 Redis高并发场景
基于GBase8s和Calcite的多数据源查询
Redis——底层数据结构原理
上一篇文章      下一篇文章      查看所有文章
加:2022-04-01 23:28:20  更:2022-04-01 23:31:03 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 -2025/1/16 14:56:00-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码