Apache Flink 集成 Hadoop InputFormat 使用指南:基于 flink-hadoop-compatibility 模块的实践与源码解析

发布时间:2026/9/20 5:05:01

Apache Flink 集成 Hadoop InputFormat 使用指南:基于 flink-hadoop-compatibility 模块的实践与源码解析 Apache Flink 集成 Hadoop InputFormat 使用指南基于 flink-hadoop-compatibility 模块的实践与源码解析【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink本文以 Apache Flink 仓库中的 Hadoop formats 官方文档 为主体系统讲解如何在 Flink DataStream 作业中复用 Hadoop 生态的InputFormat含旧版mapred与新版mapreduce两套 API。读者将掌握flink-hadoop-compatibility模块的 Maven 依赖配置、HadoopInputs工具类的核心方法readHadoopFile/createHadoopInput/readSequenceFile、Java 与 Scala 两种 API 的完整用法并从源码层面理解 Flink 包装 Hadoop InputFormat 的分片、读取、序列化与凭证传递机制从而将 HDFS 文本文件、SequenceFile 以及第三方 Hadoop InputFormat 无缝接入 Flink 作业。一、Hadoop 兼容模块与项目配置对 Hadoop 的支持位于flink-hadoop-compatibilityMaven 模块中该模块同时提供 Java 与 Scala 两套 API。在模块源码树中可以清晰看到其组织方式Java API 核心位于 flink-hadoop-compatibility/src/main/java/org/apache/flink 下包含工具类org.apache.flink.hadoopcompatibility.HadoopInputs、mapred与mapreduce两套 InputFormat/OutputFormat 包装实现以及Writable类型的序列化支持Scala API 位于 flink-hadoop-compatibility/src/main/scala/org/apache/flink 下对应org.apache.flink.hadoopcompatibility.scala.HadoopInputs与org.apache.flink.api.scala.hadoop包。1.1 添加 Maven 依赖在使用 Hadoop InputFormat 的 Flink 项目pom.xml中添加如下依赖dependency groupIdorg.apache.flink/groupId artifactIdflink-hadoop-compatibility{{ scala_version }}/artifactId version{{ version }}/version /dependency其中{{ scala_version }}为构建所用的 Scala 二进制版本后缀如_2.12{{ version }}为当前 Flink 版本。事实上在 flink-hadoop-compatibility/pom.xml 中该模块的 artifactId 正是以${scala.binary.version}动态拼接的artifactIdflink-hadoop-compatibility_${scala.binary.version}/artifactId且模块自身的依赖hadoop-common、hadoop-mapreduce-client-core均声明为provided作用域说明 Hadoop 相关类库需要由运行环境或用户显式提供而非随该模块打包分发。1.2 本地运行时的 Hadoop 客户端依赖如果你想在本地运行 Flink 应用例如在 IDE 中还需要按照如下所示将hadoop-client依赖也添加到pom.xmldependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version2.10.2/version scopeprovided/scope /dependency这里给出几点实践建议hadoop-client是一个聚合依赖会传递引入hadoop-common、hadoop-hdfs、hadoop-mapreduce-client-core等常用模块适合本地开发调试使用provided作用域可以避免把 Hadoop 类库打进 Fat JAR 与 Flink 自身依赖冲突部署到集群时由集群侧的 Hadoop 环境HADOOP_CLASSPATH、FLINK_HADOOP_CLASSPATH等提供具体 Hadoop 版本应与目标集群版本保持一致文档给出的2.10.2是经 Flink 兼容性验证的参考版本实际使用时请以运行环境为准。二、HadoopInputs 工具类三种工厂方法与两套 API在 Flink 中使用 HadoopInputFormat必须首先使用HadoopInputs工具类进行包装。HadoopInputs是一个final工具类其 Java 实现位于 HadoopInputs.javaScala 实现位于 scala/HadoopInputs.scala。两者提供的工厂方法一一对应。2.1 readHadoopFile包装 FileInputFormat 系输入格式readHadoopFile用于包装从org.apache.hadoop.mapred.FileInputFormat旧 API或org.apache.hadoop.mapreduce.lib.input.FileInputFormat新 API派生的 Input Format。以 Javamapred版本为例public static K, V HadoopInputFormatK, V readHadoopFile( org.apache.hadoop.mapred.FileInputFormatK, V mapredInputFormat, ClassK key, ClassV value, String inputPath, JobConf job)从源码可以看到它的内部逻辑是先调用 Hadoop 的FileInputFormat.addInputPath(job, new Path(inputPath))把输入路径写入JobConf再委托给createHadoopInput完成包装。此外还提供了省略JobConf的重载内部自动new JobConf()以及readSequenceFile(ClassK key, ClassV value, String inputPath)便捷方法——后者内部直接实例化 Hadoop 的SequenceFileInputFormat来读取 Hadoop SequenceFile无需手动构造 InputFormat 对象。2.2 createHadoopInput包装通用 InputFormatcreateHadoopInput用于包装通用的 HadoopInputFormat不限于文件类它直接把传入的 InputFormat、键值类型与JobConf或Job封装进 Flink 的HadoopInputFormat包装类public static K, V HadoopInputFormatK, V createHadoopInput( org.apache.hadoop.mapred.InputFormatK, V mapredInputFormat, ClassK key, ClassV value, JobConf job) { return new HadoopInputFormat(mapredInputFormat, key, value, job); }2.3 关键设计mapred 与 mapreduce 双 API 支持HadoopInputs对 Hadoop 两代 API 都提供了完整支持工厂方法mapred旧 APIorg.apache.hadoop.mapredmapreduce新 APIorg.apache.hadoop.mapreducereadHadoopFile接受mapred.FileInputFormat用JobConf承载配置接受mapreduce.lib.input.FileInputFormat用Job承载配置createHadoopInput返回org.apache.flink.api.java.hadoop.mapred.HadoopInputFormat返回org.apache.flink.api.java.hadoop.mapreduce.HadoopInputFormatreadSequenceFile基于mapred.SequenceFileInputFormat同左SequenceFile 读取同样位于 mapred 包两个 API 版本的包装类位于不同包路径org.apache.flink.api.java.hadoop.mapred.HadoopInputFormat与org.apache.flink.api.java.hadoop.mapreduce.HadoopInputFormat它们均继承各自包下的HadoopInputFormatBase对外行为一致产出Tuple2K, Vf0 为键、f1 为值。选择哪套 API 取决于你手中的 Hadoop InputFormat 属于哪个版本——老代码多为mapred新生态如基于 mapreduce 编写的第三方格式用mapreduce。三、使用示例从 Hadoop 文本文件创建 DataStream包装完成后生成的 FlinkInputFormat可通过StreamExecutionEnvironment#createInput批环境为ExecutionEnvironment#createInput创建数据源。生成的DataStream包含 2 元组其中第一个字段是键第二个字段是从 HadoopInputFormat接收的值。下面以 Hadoop 的KeyValueTextInputFormat为例该格式按key \t value切分文本行给出 Java 与 Scala 两种写法。3.1 Java 版本StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); KeyValueTextInputFormat textInputFormat new KeyValueTextInputFormat(); DataStreamTuple2Text, Text input env.createInput(HadoopInputs.readHadoopFile( textInputFormat, Text.class, Text.class, textPath)); // Do something with the data. [...]3.2 Scala 版本val env StreamExecutionEnvironment.getExecutionEnvironment val textInputFormat new KeyValueTextInputFormat val input: DataStream[(Text, Text)] env.createInput(HadoopInputs.readHadoopFile( textInputFormat, classOf[Text], classOf[Text], textPath)) // Do something with the data. [...]3.3 完整可运行示例mapreduce API TextInputFormat仓库中的集成测试 WordCountMapreduceITCase.java 给出了一个完整可运行的 WordCount 案例完整展示了 Hadoop InputFormat 的接入流程读取 → 转换 → 分组聚合 → 写回 Hadoop OutputFormatfinal ExecutionEnvironment env ExecutionEnvironment.getExecutionEnvironment(); // 1) 用 readHadoopFile 包装 mapreduce 版的 TextInputFormat产出 Tuple2LongWritable, Text DataSetTuple2LongWritable, Text input env.createInput( HadoopInputs.readHadoopFile( new TextInputFormat(), LongWritable.class, Text.class, textPath)); // 2) 取出 valueText转换为字符串 DataSetString text input.map(value - value.f1.toString()); // 3) 常规 Flink 处理分词、分组、求和 DataSetTuple2String, Integer counts text.flatMap(new Tokenizer()) .groupBy(0) .sum(1); // 4) 转换回 Hadoop Writable 类型通过 HadoopOutputFormat 写回 HDFS DataSetTuple2Text, LongWritable words counts.map(v - new Tuple2(new Text(v.f0), new LongWritable(v.f1))); Job job Job.getInstance(); HadoopOutputFormatText, LongWritable hadoopOutputFormat new HadoopOutputFormat(new TextOutputFormat(), job); job.getConfiguration().set(mapred.textoutputformat.separator, ); TextOutputFormat.setOutputPath(job, new Path(resultPath)); words.output(hadoopOutputFormat); env.execute(Hadoop Compat WordCount);该测试覆盖了读 Hadoop 文件 → Flink 处理 → 写回 Hadoop 文件的完整闭环是学习接入姿势的最佳参考。值得注意的细节是读取端使用LongWritable/Text作为键值类型写回端同样使用Text/LongWritable——因为 Hadoop Writable 类型的序列化与 Flink 内置类型不同需要专门支持见下文第五节。四、底层原理Flink 如何包装 Hadoop InputFormat了解包装类的内部实现有助于排查并行度、分片、性能相关问题。以 mapred 版本 HadoopInputFormatBase.java 为例其生命周期与 Flink 的RichInputFormat完全对齐构造与配置configure构造时通过HadoopUtils.mergeHadoopConf(job)将 Hadoop 配置合并进 Flink 环境并用ReflectionUtils.setConf把JobConf注入 InputFormatconfigure()阶段则对实现Configurable或JobConfigurable接口的 InputFormat 再次注入配置。需要注意Flink 并行度远超 Hadoop 的多进程模型多个 InputFormat 实例可能运行在同一个 JVM 的不同线程中因此基类专门用OPEN_MUTEX、CONFIGURE_MUTEX、CLOSE_MUTEX三个静态互斥锁串行化open/configure/close调用避免某些依赖 JVM 隔离的 Hadoop 实现产生并发问题。分片createInputSplitsmapred 版本调用mapredInputFormat.getSplits(jobConf, minNumSplits)生成 Hadoop 原生InputSplit数组再用 HadoopInputSplit 包装为 Flink 的InputSplit并通过LocatableInputSplitAssigner支持基于位置的本地性调度。mapreduce 版本在 mapreduce/HadoopInputFormatBase.java 中还会设置mapreduce.input.fileinputformat.split.minsize配置后再调用getSplits。打开与读取open / nextRecordopen(split)中调用mapredInputFormat.getRecordReader(split, jobConf, new HadoopDummyReporter())创建RecordReadermapreduce 版本为createRecordReaderinitialize配合TaskAttemptContextImpl其中HadoopDummyReporter是一个空实现因为 Flink 有自己独立的进度与计数器体系随后nextRecord把record.f0 key; record.f1 value填充进Tuple2。序列化与反序列化writeObject / readObject包装类实现了自定义的 Java 序列化逻辑writeObject只写入 InputFormat 的类名、键值类名与序列化后的JobConfreadObject在反序列化时通过Class.forName反射重新实例化 InputFormat并把JobConf中携带的Credentials与当前UserGroupInformation的凭证合并从而支持 Kerberos 等安全认证场景基类 HadoopInputFormatCommonBase.java 专门负责凭证的读写与传递。这也解释了为什么 Hadoop 的InputFormat实现类与键值类必须出现在任务执行的 classpath 上——反序列化依赖类名反射。统计信息getStatistics仅当包装的是FileInputFormat时才会枚举输入路径下的文件含目录递归计算总大小与最新修改时间并生成FileBaseStatistics供 Flink 优化器评估数据规模。五、Writable 类型的序列化支持HadoopInputFormat产出的键值通常是org.apache.hadoop.io.Writable的子类如Text、LongWritable、IntWritable、NullWritable等它们并不实现 Java 的Serializable无法直接由 Flink 默认序列化器处理。为此该模块提供了专门支持WritableTypeInfo.java 用于识别Writable类型WritableSerializer.java 实现了TypeSerializerT extends Writable其序列化策略非常巧妙利用 Writable 自描述的write(DataOutput)/readFields(DataInput)接口完成二进制序列化serialize调用record.write(target)deserialize调用reuse.readFields(source)复用对象而对象复制copy则通过 Kryo 完成并注册了typeClass。同时实现了TypeSerializerSnapshotWritableSerializerSnapshot支持作业恢复时对序列化器配置的校验与兼容保证从保存点Savepoint恢复作业的稳定性。正是由于该序列化器的存在Tuple2Text, Text、Tuple2LongWritable, Text这类包含 Writable 类型的数据流才能被 Flink 正常持久化、分发与恢复。六、使用注意事项与最佳实践结合文档与仓库源码总结以下实践要点API 版本匹配readHadoopFile的mapred重载接收org.apache.hadoop.mapred.FileInputFormat与JobConfmapreduce重载接收org.apache.hadoop.mapreduce.lib.input.FileInputFormat与Job两者不能混用导入包时务必留意。键值类型必须给出所有工厂方法都要求显式传入ClassK与ClassV这些类型既用于getProducedType()推导TypeInformationTuple2K, V见 mapred/HadoopInputFormat.java 的getProducedType也用于反序列化时反射重建类写错会导致运行期类型错误。依赖作用域与版本flink-hadoop-compatibility对 Hadoop 依赖为provided本地 IDE 运行时务必按第一节添加hadoop-client参考版本2.10.2部署到 YARN/Kubernetes 集群时应依赖集群自带的 Hadoop 环境避免版本冲突。集群执行时还应通过HADOOP_CLASSPATH等机制确保 Hadoop 类库在 TaskManager 的 classpath 中。路径与安全readHadoopFile会把输入路径写入JobConf/Job支持 HDFS 路径如hdfs://namenode:8020/path与本地路径若集群开启了 Kerberos模块会在序列化与createInputSplits阶段合并Credentials凭证保证安全认证正常流转。并行度与分片Flink 的并行度由createInputSplits产出的 Hadoop Split 数量决定每个 Split 分配给一个并行子任务若要控制并发可在调用createInput后用setParallelism约束但不能超过 Split 总数。操作系统限制仓库测试类WordCountMapreduceITCase中通过assumeThat(OperatingSystem.isWindows()).isFalse()跳过 Windows 上的执行原因指向 Hadoop 在 Windows 平台的历史兼容问题FLINK-5164开发调试时建议优先在 Linux/macOS 环境验证。Scala API 已弃用Scala 版本的 HadoopInputs.scala 自 Flink 1.18.0 起被标记为deprecated依据 FLIP-265所有 Flink Scala API 将被逐步移除新项目建议直接使用 Java API 编写 DataStream 作业。七、小结通过flink-hadoop-compatibility模块Flink 可以无缝复用 Hadoop 生态中数量庞大的InputFormat实现只需三步——添加 Maven 依赖、用HadoopInputs.readHadoopFile/createHadoopInput包装、通过env.createInput创建数据源。生成的DataStreamTuple2K, V可以直接参与 Flink 的转换、窗口、聚合与状态管理等操作底层由HadoopInputFormatBase完成分片映射、RecordReader 适配、Writable 序列化与凭证传递。无论你是想读取 HDFS 文本、SequenceFile还是集成某个基于 Hadoop InputFormat 的第三方数据源本文给出的配置、示例与源码分析都已覆盖完整的接入路径。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/20 5:05:01

DeepSeek API 401 报错排查清单:从 Key 到代理的完整链路

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 5:00:01

AssetRipper完整指南:快速跑通Unity资源提取

AssetRipper完整指南:快速跑通Unity资源提取 【免费下载链接】AssetRipper GUI application to analyze game files 项目地址: https://gitcode.com/GitHub_Trending/as/AssetRipper AssetRipper 是一款免费开源的 Unity资源提取工具:它解析 .ass…

2026/9/20 6:15:04

龙珠Z第193集:神龙升级与角色成长解析

1. 龙珠Z第193集深度解析:愿望与抉择的哲学《龙珠Z》第193集展现了丹迪使用改造后的龙珠召唤出升级版神龙的关键情节。这一集不仅推动了剧情发展,更通过角色间的互动揭示了深刻的主题内涵。新神龙能够实现两个愿望的能力设定,为后续故事埋下了…

2026/9/20 6:15:04

路由与导航系统:核心架构设计与工程实践

1. 路由与导航系统概述在移动应用和Web开发领域,路由与导航系统就像城市交通的GPS导航,它决定了用户如何在不同界面间跳转流转。我经历过多个大型项目后深刻体会到:优秀的导航设计能让用户像走在熟悉的街道上,而糟糕的实现则会让应…

2026/9/20 6:15:04

TensorRT部署实战:YOLO转ONNX到推理加速的五大避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/20 6:15:04

AI如何优化论文格式审查与投稿效率

1. 项目背景与痛点解析学术论文发表是每个研究者必经的"修罗场"。从选题构思到最终见刊,平均需要经历17.3次修改(Nature指数统计),其中约68%的投稿因格式合规性问题被直接拒稿。更令人焦虑的是,Elsevier旗下…

2026/9/20 6:15:04

AI原生研发组织转型实践:从辅助工具到流程重构的深度复盘

最近有大半年时间,我基本没怎么在公开场合系统聊过我们团队在AI研发组织上的做法,不是藏着掖着,是确实一直在试错和调整。各个群里被问得多了,索性把这几个月的一些探索和实践整理一下。这算是一篇比较完整的复盘,不吹…

2026/9/20 6:10:04

OpenResearch:多AI编程工具协作的上下文管理与复现工作流

1. 从"OpenResearch"这个名字说起:它到底想解决什么问题第一次看到"OpenResearch"这个标题,加上旁边一串 Claude Code、Codex、OpenCode、Cursor 的热搜词,我大概能猜到它想干的事:把当下最火的几个 AI 编程工…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 4:54:47

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/20 5:01:23

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/20 5:09:33

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码