用Flink对keyed state 状态的代码进行设置TTL(过期时间)和hdfs的checkpoints恢复和配置Flink-conf.yaml
依赖dependencygroupIdorg.apache.flink/groupIdartifactIdflink-core/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-clients_2.11/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-scala_2.11/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-streaming-scala_2.11/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.slf4j/groupIdartifactIdslf4j-log4j12/artifactIdversion1.7.26/version/dependency!--添加HDFS支持--dependencygroupIdorg.apache.hadoop/groupIdartifactIdhadoop-hdfs/artifactIdversion2.9.2/version/dependencydependencygroupIdorg.apache.hadoop/groupIdartifactIdhadoop-common/artifactIdversion2.9.2/version/dependency!--添加Kafka支持--dependencygroupIdorg.apache.flink/groupIdartifactIdflink-connector-kafka_2.11/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-json/artifactIdversion1.8.1/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-connector-filesystem_2.11/artifactIdversion1.8.1/version/dependency!--写到Redis中--dependencygroupIdorg.apache.bahir/groupIdartifactIdflink-connector-redis_2.11/artifactIdversion1.0/version/dependency/dependenciesbuildplugins!--在执行package时候将scala源码编译进jar--plugingroupIdnet.alchim31.maven/groupIdartifactIdscala-maven-plugin/artifactIdversion4.0.1/versionexecutionsexecutionidscala-compile-first/idphaseprocess-resources/phasegoalsgoaladd-source/goalgoalcompile/goal/goals/execution/executions/plugin/plugins/build代码packagecom.baizhi.demo04importorg.apache.flink.streaming.api.CheckpointingModeimportorg.apache.flink.streaming.api.environment.CheckpointConfig.ExternalizedCheckpointCleanupimportorg.apache.flink.streaming.api.scala.{DataStream,StreamExecutionEnvironment,_}object Flink{defmain(args:Array[String]):Unit{val envStreamExecutionEnvironment.getExecutionEnvironment//开启checkpointenv.enableCheckpointing(7000,CheckpointingMode.EXACTLY_ONCE)//checkpoint必须在2s内完成如果完成不了终止env.getCheckpointConfig.setCheckpointTimeout(4000)//距离上一次的checkpoint完成之后需要等5s 之后再开启下一次的checkpointenv.getCheckpointConfig.setMinPauseBetweenCheckpoints(5000)env.getCheckpointConfig.setMaxConcurrentCheckpoints(1)//只开启一个checkpoint线程//在退出应用时候不删除checkpoint数据env.getCheckpointConfig.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION)//必须保证任务可以从checkpoint恢复恢复不成功任务失败env.getCheckpointConfig.setFailOnCheckpointingErrors(true)val data:DataStream[String]env.socketTextStream(Flink,9999)data.flatMap(lineline.split(\\s)).map((_,1)).keyBy(0).map(newCountMapFunction).print()env.execute()}}packagecom.baizhi.demo04importorg.apache.flink.api.common.functions.RichMapFunctionimportorg.apache.flink.api.common.state.{StateTtlConfig,ValueState,ValueStateDescriptor}importorg.apache.flink.api.common.time.Timeimportorg.apache.flink.configuration.Configurationimportorg.apache.flink.streaming.api.scala._classCountMapFunctionextendsRichMapFunction[(String,Int),(String,Int)]{var state:ValueState[Int]_ override defmap(value:(String,Int)):(String,Int){var historystate.value()if(historynull){history0}state.update(historyvalue._2)(value._1,historyvalue._2)}override defopen(parameters:Configuration):Unit{val decnewValueStateDescriptor[Int](count,createTypeInformation[Int])//1.创建TTLConfigval ttlConfigStateTtlConfig.newBuilder(Time.seconds(80))//这是state存活时间10s.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)//设置过期时间更新方式.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)//永远不要返回过期的状态//.cleanupInRocksdbCompactFilter(1000)//处理完1000个状态查询时候会启用一次CompactFilter.build//2.开启TTLdec.enableTimeToLive(ttlConfig)val contextgetRuntimeContext statecontext.getState(dec)}}配置Flink的环境配置Flink-conf.yaml## Fault tolerance and checkpointing## The backend that will be used to store operator state checkpoints if# checkpointing is enabled.## Supported backends are jobmanager, filesystem, rocksdb, or the# class-name-of-factory.#state.backend:rocksdb# Directory for checkpoints filesystem, when using any of the default bundled# state backends.#state.checkpoints.dir:hdfs:///flink-checkpoints# Default target directory for savepoints, optional.#state.savepoints.dir:hdfs:///flink-savepoints# Flag to enable/disable incremental checkpoints for backends that# support incremental checkpoints (like the RocksDB state backend).#state.backend.incremental:truestate.backend.rocksdb.ttl.compaction.filter.enabled:true## HistoryServer## The HistoryServer is started and stopped via bin/historyserver.sh (start|stop)# Directory to upload completed jobs to. Add this directory to the list of# monitored directories of the HistoryServer as well (see below).jobmanager.archive.fs.dir:hdfs:///completed-jobs/# The address under which the web-based HistoryServer listens.historyserver.web.address:CentOS# The port under which the web-based HistoryServer listens.historyserver.web.port:8082# Comma separated list of directories to monitor for completed jobs.historyserver.archive.fs.dir:hdfs:///completed-jobs/# Interval in milliseconds for refreshing the monitored directories.historyserver.archive.fs.refresh-interval:10000将hadoop_classpath配置到环境变量中export HADOOP_CLASSPATHhadoop classpath

相关新闻

单细胞测序技术解析:CIMA免疫细胞图谱与应用

单细胞测序技术解析:CIMA免疫细胞图谱与应用

1. 项目背景与意义 CIMA(Cell-level Immune Map Atlas)千万级免疫细胞图谱的发布,标志着单细胞测序技术在免疫学研究领域取得重大突破。这个项目通过对上千万个免疫细胞进行高通量测序和系统性分析,构建了迄今为止最全面的免疫细胞…

2026/7/28 13:52:35 阅读更多 →
Hashcat实战:五种方法破解NTLM哈希与掩码规则速查

Hashcat实战:五种方法破解NTLM哈希与掩码规则速查

1. 项目概述:当NTLM哈希遇上Hashcat 在渗透测试和红队评估的日常工作中,我们经常会遇到从目标系统中提取出的用户凭证哈希。其中,NTLM(NT LAN Manager)哈希是Windows环境中最常见、也最核心的一种凭证存储形式。它不像…

2026/7/28 13:52:35 阅读更多 →
C#中的Nutshell函数式编程

C#中的Nutshell函数式编程

目录 介绍 函数编程定义 函数属性 纯度 头等函数 闭包的概念 成为函数式 函数式实用程序 纯度重要性 头等的重要性 函数编程和面向对象编程 集成函数编程 总结 下载源码 - 7.9 KB 介绍 如今,函数性编程正在流行。我们应该问自己有两个问题&#xff1a…

2026/7/28 13:52:35 阅读更多 →

最新新闻

风扇静音革命:用Fan Control为你的PC打造“隐形“散热管家

风扇静音革命:用Fan Control为你的PC打造“隐形“散热管家

风扇静音革命:用Fan Control为你的PC打造"隐形"散热管家 【免费下载链接】FanControl.Releases This is the release repository for Fan Control, a highly customizable fan controlling software for Windows. 项目地址: https://gitcode.com/GitHub…

2026/7/28 14:23:12 阅读更多 →
如何用Diablo Edit2存档修改器告别暗黑2重复刷怪,3分钟打造完美角色?

如何用Diablo Edit2存档修改器告别暗黑2重复刷怪,3分钟打造完美角色?

如何用Diablo Edit2存档修改器告别暗黑2重复刷怪,3分钟打造完美角色? 【免费下载链接】diablo_edit Diablo II Character editor. 项目地址: https://gitcode.com/gh_mirrors/di/diablo_edit 你是否厌倦了在暗黑破坏神2中花费数小时只为刷一件装备…

2026/7/28 14:23:12 阅读更多 →
成本CO

成本CO

物料成本构成 TCK07 成本构成结构 (用公司代码取elehk)TCKH3 成本构成结构和对应的列ckmlhd 成本估算号(物料工厂得到一个成本估算号)marc 取成本核算批量ckmlprkeph 标准成本 (成本估算号年度期间 得到一条数据&#…

2026/7/28 14:23:12 阅读更多 →
RF射频电路PCB布局布线核心技巧与实战经验

RF射频电路PCB布局布线核心技巧与实战经验

1. RF射频信号布局布线的重要性与挑战在电子设计领域,RF射频电路一直被视为"黑魔法"般的存在。作为一名经历过无数次改板的PCB工程师,我深刻理解射频设计中的每一个微小失误都可能导致整个项目推倒重来。与低频数字电路不同,射频信…

2026/7/28 14:23:12 阅读更多 →
AI编程助手横向评测:ChatGPT、Gemini与Claude在游戏复刻项目中的实战表现

AI编程助手横向评测:ChatGPT、Gemini与Claude在游戏复刻项目中的实战表现

作为一名游戏开发者,你是否曾想过,如果让当前最顶尖的AI编程助手来“复刻”一款经典游戏,它们各自的表现会如何?是ChatGPT的代码逻辑更严谨,Gemini的创意更贴合需求,还是Claude的工程化能力更强&#xff1f…

2026/7/28 14:23:12 阅读更多 →
快手视频批量下载神器:3分钟学会无水印保存技巧

快手视频批量下载神器:3分钟学会无水印保存技巧

快手视频批量下载神器:3分钟学会无水印保存技巧 【免费下载链接】KS-Downloader 快手(KuaiShou)作品视频/图片下载工具 项目地址: https://gitcode.com/gh_mirrors/ks/KS-Downloader 还在为无法保存喜欢的快手视频而烦恼吗&#xff1f…

2026/7/28 14:22:12 阅读更多 →

日新闻

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生 【免费下载链接】OmenSuperHub Control Omen laptop performance, fan speeds, and keyboard lighting, and unlock power limits. 项目地址: https://gitcode.com/gh_mirrors/om/OmenSuperHub 你是否也曾为官方Om…

2026/7/28 0:00:43 阅读更多 →
RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

做 RAG 的人应该都踩过这个致命的坑:把几百页的财报、法规、技术手册扔给向量库,问一个具体问题,搜出来的全是沾边但没用的内容 —— 关键信息要么被硬切块拆碎了,要么藏在几十条结果的最下面。语义相似≠真正相关,这个…

2026/7/28 0:00:43 阅读更多 →
抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

2026年做短视频运营,从抖音上扒文案早就不是偷偷抄笔记的事了。我刚开始做内容的时候,每天刷半小时抖音,手动把爆款视频的口播敲进备忘录,一条2分钟的视频得花十来分钟,碰到语速快的还要反复回听。后来试了一圈工具&am…

2026/7/28 0:00:43 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/28 12:04:22 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/28 8:29:16 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/28 5:03:42 阅读更多 →

月新闻