ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

HDFS多文件压缩合并下载:GzipCodec与IOUtils实战解析

HDFS多文件压缩合并下载:GzipCodec与IOUtils实战解析 简介这份云计算实验报告聚焦Hadoop IO文件读写与Gzip压缩编程面向高校计算机专业学生及云计算自学者解决云端多个文件合并压缩下载到本地的实操问题。资源以pdf格式呈现共1个文件大小573KB内容完整可对照实验流程逐步学习。已有307人学习适合正在完成实验作业或想巩固Hadoop IO知识的人参考。报告基于Linux与Eclipse环境详细记录了改写GetMerge程序实现多文件压缩下载的过程包含实验目的、实验要求、编程步骤、核心代码及结果总结并通过Gzip压缩方法将多个云端文件合并为一个压缩文件。作为一份93分的实验报告它不仅能帮助读者快速掌握Hadoop中文件读写、压缩流创建与数据复制等关键操作还能从输入输出流配置、压缩格式设置等细节中理解大数据存储优化思路为深入学习云计算打下基础。1. 从 GetMerge 到 GetMerge1HDFS 多文件下载为什么要卡在压缩这一步实验五和实验四的差距从代码量看只有三行创建 GzipCodec、包装 CompressionOutputStream、把 copyBytes 的输出流替换成压缩流。但把这三行放进整条数据链路看会得到一个反直觉的结论——GetMerge1 并不是一个 MapReduce 作业它没有 Mapper 也没有 Reducer不需要向 YARN 申请容器它本质上是跑在 JVM 里的 HDFS 客户端程序完成的是把云端多个文件按顺序读出、压缩后写入本地 .gz 文件这件并不分布式的事。实验之所以把压缩放在合并下载这一步是因为跨节点搬运数据时带宽成本远高于本地 CPU 压缩成本。一个 Merger.gz 从 HDFS 拉到本地体积通常是原始日志文件的 20%35%传输时间明显缩短。如果是按流量计费的跨机房下载这个差距直接反映在账单上。所以 GetMerge1 的核心不是 MapReduce 编程模型而是 Hadoop IO 里压缩编解码器的工作方式。2. Hadoop IO 读写链路上的 Compress 层CompressionCodec 接口与 GzipCodec 行为2.1 HDFS 读写链路中压缩实际发生在哪一层HDFS 读路径的常规流程是客户端通过 DistributedFileSystem.open 拿到 DistributedFileSystem 的读入口内部经 DFSClient 定位 Block 所在的 DataNode再建立流式读取返回 FSDataInputStream。写路径类似通过 DFSOutputStream 把数据切块后发往 DataNode 副本管线。这两条链路里都不存在压缩逻辑HDFS 自身只负责块存储和副本管理IO 压缩是应用侧通过装饰流方式额外挂上去的。GetMerge1 的数据流可以这样看FSDataInputStream远端 HDFS 文件→ IOUtils.copyBytes 内部 4096 字节缓冲 → CompressionOutputStream → FileOutputStream本地 Merger.gz。关键点是 CompressionOutputStream 包在本地文件流外面所有来自远端文件的字节在落盘之前先经过 java.util.zip.Deflater 压缩。每写入一个文件压缩器继续累积字典数据因此多个高相似度的日志文件合并写入同一个 Gzip 流时压缩率通常高于分别压缩再拼接文件的总和。这里还有一个容易误判的地方压缩后的 Merger.gz 是单条 gzip 流解压时只需要一个解压器从头读到尾。如果实验四的做法是先把多个文件拼成一个大文件再单独压缩那么压缩合流的顺序会影响最终结果而 GetMerge1 是边读边压远端文件内容直接进入同一个压缩上下文得到的 .gz 文件内部不存在文件边界标记。2.2 CompressionCodec 与 ReflectionUtils.newInstance 的构造方式Hadoop 把压缩算法统一抽象为 CompressionCodec 接口核心方法只有两个createOutputStream(OutputStream out) 把现有输出流包装成压缩输出流createInputStream(InputStream in) 把压缩输入流解包为原始数据流。上层代码只需要面向这个接口编程不关心底层是 gzip、bzip2 还是 snappy替换算法时只换实现类即可。代码里用的 ReflectionUtils.newInstance(GzipCodec.class, conf) 是 Hadoop 组件最常见的实例化方式而不是直接 new GzipCodec()。原因是 ReflectionUtils 会先处理 Configuration 相关的依赖注入并适配 Hadoop 服务加载机制确保在不同 Hadoop 版本下 GzipCodec 内部引用的压缩器实现能被正确初始化。实验里如果换成直接 new在部分 CDH 或自编译版本下会抛出 NoSuchMethodError这属于典型的版本兼容问题。GzipCodec 位于 org.apache.hadoop.io.compress 包底层是 JDK 自带的 Deflater/Inflater不需要额外 native 库这是它适合课程实验的根本原因。下表是常见编解码器的对比生产环境选择时基本也按这个维度评估实现类压缩率文本日志压缩速度native 依赖常见场景GzipCodec中中无跨集群下载、日志归档BZip2Codec高慢是冷数据压缩存储SnappyCodec低快是MapReduce/Tez 中间结果Lz4Codec低很快是高频临时数据落盘实验要求用 Gzip 而不是 Snappy除了题目约束外更实际的考虑是 native 库问题。SnappyCodec 在 Windows 开发机上经常出现 hadoop-snappy 的 native 库加载失败而 GzipCodec 纯 JVM 实现伪分布式、Eclipse 直接运行都不会有环境障碍。理解这一层的选型逻辑比背下三行代码更能应对后续问题。3. GetMerge1 工程落地Eclipse 依赖、多文件遍历与压缩输出流绑定3.1 MapReduce 工程创建与依赖范围Eclipse 里新建 MapReduce Project 时实际生成的是带 Hadoop 类库引用的普通 Java 工程。不必把它理解为必须提交到 YARN 运行的任务工程GetMerge1 的 main 方法在本地 JVM 跑起来后以客户端身份访问 HDFS整个过程不产生作业 ID。工程导入的 jar 包以 hadoop-client 为主。Maven 工程里最简配置为dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version /dependencyhadoop-client 会传递引入 hadoop-common、hadoop-hdfs、hadoop-mapreduce-client-core 等模块覆盖 GzipCodec、ReflectionUtils、FileSystem 所在包。非 Maven 工程在 Eclipse Build Path 里手动添加 hadoop-common、hadoop-hdfs 和 hadoop-client jar 也可以但注意只引一份避免 guava、protobuf 多个版本同时出现在 classpath 里否则运行时会先报加载冲突而不是业务错误。3.2 多文件枚举、压缩流绑定与文件循环写入GetMerge1 的完整代码围绕三个动作展开获取远端目录文件列表、构造 Gzip 压缩流、逐个文件循环拷贝。下面是可直接运行的实现package lab5; import java.io.FileOutputStream; import java.io.OutputStream; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FSDataInputStream; import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IOUtils; import org.apache.hadoop.io.compress.CompressionCodec; import org.apache.hadoop.io.compress.CompressionOutputStream; import org.apache.hadoop.io.compress.GzipCodec; import org.apache.hadoop.util.ReflectionUtils; public class GetMerge1 { public static void main(String[] args) throws Exception { if (args.length 2) { System.err.println(Usage: GetMerge1 hdfs-src-dir local-gz-file); System.exit(1); } Configuration conf new Configuration(); Path srcDir new Path(args[0]); String localFile args[1]; FileSystem fs srcDir.getFileSystem(conf); FileStatus[] stats fs.listStatus(srcDir); // 实验要求云端文件需要超过 2 个 if (stats.length 3) { System.err.println(warning: files less than 3, got stats.length); } // 创建压缩器 codec CompressionCodec codec (CompressionCodec) ReflectionUtils.newInstance(GzipCodec.class, conf); OutputStream out new FileOutputStream(localFile); // 创建输出压缩流 CompressionOutputStream outGzip codec.createOutputStream(out); for (FileStatus status : stats) { if (status.isDirectory()) { continue; } FSDataInputStream in fs.open(status.getPath()); // 远端输入流 in - 压缩输出流 outGzip缓冲 4096 字节 IOUtils.copyBytes(in, outGzip, 4096, false); in.close(); } outGzip.finish(); outGzip.close(); fs.close(); System.out.println(merged gzip saved to localFile); } }FileSystem fs srcDir.getFileSystem(conf) 这一行是从 Path 的 URI scheme 推断文件系统类型hdfs:// 开头连 HDFSfile:// 开头连本地因此同样一段代码可以在本地文件系统联调。listStatus 返回的数组同时包含文件和目录程序里必须过滤 isDirectory()否则 fs.open 一个目录路径时抛 FileNotFoundException而实验里云端路径下很容易混入 _SUCCESS 这样的空文件或子目录。ReflectionUtils.newInstance(GzipCodec.class, conf) 返回的 codec 是关键入口createOutputStream 把本地 OutPutStream 包装成 CompressionOutputStream。后续所有 copyBytes 都写向它数据经 Deflater 压缩后进入本地文件。IOUtils.copyBytes(in, outGzip, 4096, false) 的四个参数依次为输入流、输出流、缓冲字节数、是否在拷贝完成后关闭输出流。这里第四个参数传 false 是刻意的循环里还要继续复用 outGzip不能拷贝一个文件夹就关一次。循环结束后调用 outGzip.finish() 而非直接 close这一点非常关键。CompressionOutputStream 内部有压缩器状态finish 会把压缩器剩余数据刷出并写出 gzip 格式要求的 CRC-32 校验值和 ISIZE 字段如果跳过 finish 直接关闭底层流得到的 Merger.gz 在解压时大概率报 Incorrect header check 或在尾部截断。程序最后再 close 压缩流和 FileSystem保证文件句柄释放。3.3 4096 缓冲参数的含义与调整思路IOUtils.copyBytes 的第三个参数 4096 与 Hadoop 配置项 io.file.buffer.size 的默认值一致它决定每次 read/write 的字节数组大小。这个参数只影响系统调用次数和内存占用不影响压缩质量和最终文件内容完整性。4096 字节意味着每次最多读入 4KB 数据交给 Deflater 处理对本地实验完全够用。如果换成生产环境下载大文件这个缓冲可以调到 65536 或 131072减少 JVM 与操作系统之间的上下文切换次数吞吐量通常有明显提升。但调大 buffer 前要先确认这是网络受限场景还是 CPU 受限场景跨机房带宽只有几十 Mbps 时buffer 翻倍收益不大瓶颈在网络内网 10Gbps 环境下大 buffer 配合压缩才能把磁盘和 CPU 都压满。4. 伪分布式运行与 Merger.gz 完整性校验安装配置之后的第一次下载排错4.1 准备三个以上 HDFS 文件并提交运行伪分布式环境搭建完成、HDFS 正常启动后先准备实验要求的云端数据。下面命令创建一个输入目录并写入三个文本文件start-dfs.sh hdfs dfs -mkdir -p /lab5/input echo 2024-06-01,page_a,120 | hdfs dfs -put - /lab5/input/log1.txt echo 2024-06-02,page_b,89 | hdfs dfs -put - /lab5/input/log2.txt echo 2024-06-03,page_c,233 | hdfs dfs -put - /lab5/input/log3.txt hdfs dfs -ls /lab5/inputstart-dfs.sh 是伪分布式搭建后的标准启停命令它依次拉起 NameNode、DataNode 和 SecondaryNameNode。若此前做过 Hadoop 集群搭建启停命令不变差别只在于多台节点时需要配合 ZKFC 实现 NameNode 自动故障转移而实验单机环境下不涉及这部分逻辑。hdfs dfs -put - 从标准输入写入 HDFS用于快速生成测试文件最方便。接着把编译好的 GetMerge1 打包提交hadoop jar hadoop-io-lab5.jar lab5.GetMerge1 /lab5/input /tmp/Merger.gzhadoop jar 命令会启动一个 JVM 进程执行 main 方法它访问 HDFS 的 NameNode 获取文件元数据再直连 DataNode 读取块数据。整个过程与 MapReduce 提交作业不同控制台不会出现 job 运行进度条程序里的 System.out 直接打印在终端。4.2 用本地命令验证 Merger.gz 是否完整运行结束后首先确认输出文件确实是 gzip 格式再做完整性和内容一致性校验file /tmp/Merger.gz gunzip -t /tmp/Merger.gz zcat /tmp/Merger.gz | wc -l hdfs dfs -text /lab5/input/log1.txt /lab5/input/log2.txt /lab5/input/log3.txt | wc -lfile 命令识别文件头gzip 文件头固定为 1f 8bgunzip -t 解压到内存并校验 CRC-32返回 silent 表示完整性通过。wc -l 统计解压后的行数与 hdfs dfs -text 直接输出源文件内容统计的行数一致时说明三个云端文件全部写入 Merger.gz没有丢数据。这里有一个与 hadoop fs -getmerge 的行为差异值得注意GetMerge1 是纯字节流拼接不会在文件之间插入换行符。如果某个源文件最后一行没有 \n合并解压后最后一行会和下一个文件的首行粘连。hadoop fs -getmerge 不带参数时也是直接拼接加 -nl 参数才插入换行。课程实验的数据通常每个文件都以换行结尾不易暴露这个问题但处理真实日志时要在合并前检查每个文件末尾是否有换行必要时在循环里补写一个 \n。4.3 编译、运行与结果层面的典型报错报错现象可能原因处理思路ClassNotFoundException: org.apache.hadoop.io.compress.GzipCodec工程缺少 hadoop-common 依赖或运行时 classpath 不完整用 hadoop jar 提交确认 jar 包包含依赖Eclipse 直接运行时检查 Build PathNoSuchMethodError 或 AbstractMethodErrorHadoop 类库版本混用比如 Apache 3.x 配了 CDH 的 client jar全工程统一 hadoop-client 版本清除重复 jarIncorrect header check / EOFExceptionoutGzip 未调用 finish 就关闭底层文件流gzip 尾部缺失确保循环结束后调用 finish() 再 close()FileNotFoundException: 路径是目录listStatus 后未过滤 isDirectory循环里加 status.isDirectory() 判断Merger.gz 能解压但行数少于源文件循环内参数传错或某个文件打开失败后异常终止打印每个 status.getPath()逐文件确认最后一类问题排查时一个实用做法是先在程序中打印 stats.length 和每个文件的长度与 hdfs dfs -ls 输出对比。如果 listStatus 返回的文件数比 ls 结果少多半是目录权限或符号链接问题如果文件数一致但内容缺失问题基本在 open/copyBytes 之间的异常处理检查是否有文件被误跳过了。5. 三个参数从实验走向生产gzip 级别、copy 缓冲与编解码器切换5.1 压缩级别io.compression.codec.gzip.levelGzipCodec 底层 Deflater 的压缩级别由配置项 io.compression.codec.gzip.level 控制范围 1 到 9默认与 JDK 一致为 6。级别越高压缩率越大但 CPU 耗时接近线性增长。在实验规模下看不出差别但下载几百 GB 日志时9 级让压缩时间翻倍体积可能只多省 3%5%。设置压缩级别的代码必须放在 ReflectionUtils.newInstance 之前因为 codec 构造阶段读取配置并初始化 Deflaterconf.setInt(io.compression.codec.gzip.level, 9); CompressionCodec codec (CompressionCodec) ReflectionUtils.newInstance(GzipCodec.class, conf);对文本日志级别 6 通常已是性价比最高的位置级别 1 适合先快速出临时文件再统一归档的场景。判断当前压缩是否拖累 IO 性能下降用 time 命令对比不同级别下同一批文件的总耗时即可。5.2 copy 缓冲与合并顺序的工程化调整IOUtils.copyBytes 的 4096 缓冲在下载大文件时可以调整为 65536同时用 Comparator 对 FileStatus 排序让合并顺序可控。实验里 listStatus 默认返回文件名字典序但真实日志系统里文件常按时间命名且有大量时间戳前缀排序逻辑必须显式写Arrays.sort(stats, Comparator.comparing(FileStatus::getModificationTime));按修改时间排序后循环写入Merger.gz 内部的日志顺序才与业务时间线一致。这行代码不影响压缩结果但直接影响下游解析程序的正确性尤其是多个小时级别的日志文件合并时顺序错乱很难在压缩层面发现。5.3 从 GzipCodec 到其他编解码器的替换边界把 GzipCodec.class 替换为 BZip2Codec.class其余代码不变即可切换算法这正是 CompressionCodec 接口带来的收益。BZip2 在文本数据上压缩率更高但压缩速度慢且 Hadoop 的 BZip2Codec 依赖 native 库某些精简环境会加载失败。SnappyCodec 相反速度快但压缩文件通常是文本日志体积的 40%50%适合 MapReduce shuffle 阶段或 Hive on Tez 的中间结果压缩。要注意的是实验中 Merger.gz 是给最终用户下载的选 Gzip 最通用因为 gzip 在所有操作系统上都有原生解压工具换成 Snappy 后用户本地反而没有现成解压器。所以生产环境里的经验是中间数据交换层用 Snappy 或 Lz4面向终端用户交付的合并文件保留 Gzip。Hive 配置 Tez 作为执行引擎时中间压缩参数 mapreduce.map.output.compress 默认也走 Snappy原理与 GetMerge1 里的 codec 机制一致。先把 GzipCodec 这个入口吃透后面接任何压缩算法都只是换一行类名的事。本文还有配套的精品资源点击获取
返回列表