恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

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

  • 首页
  • 资讯中心
  • /
  • Apache Flink 集成 Hadoop InputFormat 使用指南:基于 flink-hadoop-compatibility 模块的实践与源码解析

相关资讯

OpenToonz 2D 动画教程:5 步免费做出你的第一个弹跳球动画 2026/9/20 4:59:58
AssetRipper完整指南:快速跑通Unity资源提取 2026/9/20 4:59:58
从GitHub热榜筛选优质开源项目的5个共同点:100期观察总结 2026/9/20 4:59:58

最新资讯

油烟分离油烟机核心技术解析与选购指南
PMP认证五大过程组实战解析与项目管理黄金法则
Twitter运营实战:系统化提升内容曝光与粉丝增长
信息系统项目管理实战:从PMP到软考的核心框架解析
VLC播放器下载安装与使用全攻略:从解码到转码的实战指南
TabPFN 快速上手指南:零超参调优,1 分钟跑完表格数据分类

今日推荐

BrewUI:给Homebrew套上图形界面,让macOS软件包管理更简单
BrewUI:让Homebrew包管理变得可视化与高效
公式与文本对齐全攻略:从Word到LaTeX的实用技巧

本周热门

BrewUI:给Homebrew套上图形界面,让macOS软件包管理更简单
BrewUI:让Homebrew包管理变得可视化与高效
公式与文本对齐全攻略:从Word到LaTeX的实用技巧

本月精选

自研推理加速器Redwood:两周内实现PyTorch模型高效部署的实战教程
V4L2摄像头采集实战:从camera_client.rar到出图全流程解析
从“谁发明了钢琴键”到知识问答智能体:RAG与记忆工程实践

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

发布时间:2026/9/20 5:04:58
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),仅供参考

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号