Java读取MapReduce输出字节数组的方法及输出格式调整咨询
解决MapReduce输出字节数组的读取与直接输出问题
嘿,我来帮你搞定这个问题!你遇到的核心问题是MapReduce默认的输出格式会把BytesWritable转成字符串输出,导致读取时拿不到原始字节数组。下面分两种场景给你解决方案:
场景一:直接从现有输出文件读取字节数组
如果你的MapReduce作业已经用默认的TextOutputFormat输出了BytesWritable,那文件里其实是BytesWritable.toString()的结果——它会把字节数组按UTF-8编码转成字符串,这会导致非UTF-8字节丢失或乱码,不建议直接从这种文件恢复原始字节数组。
正确的做法是修改作业输出格式为SequenceFileOutputFormat(二进制存储Writable对象),然后用以下代码读取:
Configuration conf = new Configuration(); Path outputPath = new Path("hdfs://your-output-directory/part-r-00000"); try (SequenceFile.Reader reader = new SequenceFile.Reader(conf, SequenceFile.Reader.file(outputPath))) { // 根据文件中的键值类型实例化对象 Writable key = (Writable) ReflectionUtils.newInstance(reader.getKeyClass(), conf); BytesWritable value = (BytesWritable) ReflectionUtils.newInstance(reader.getValueClass(), conf); while (reader.next(key, value)) { // 获取原始字节数组,注意getBytes()返回的数组可能有冗余空间,用getLength()取有效长度 byte[] rawBytes = value.getBytes(); int validLength = value.getLength(); // 这里处理你的字节数组,比如: System.out.println("读取到有效字节长度:" + validLength); processRawBytes(rawBytes, 0, validLength); } } catch (IOException e) { e.printStackTrace(); }
如果实在没法修改原作业,只能读现有的文本输出文件,那你可以尝试把读取到的字符串按UTF-8转成字节数组,但仅当原始字节数组本身就是合法UTF-8字符串时才有效,代码如下:
String strFromFile = ...; // 从输出文件读取的字符串 byte[] byteArray = strFromFile.getBytes(StandardCharsets.UTF_8);
但这种方法有数据丢失风险,谨慎使用。
场景二:修改MapReduce作业直接输出原始字节数组
MapReduce的输出值必须是Writable类型,所以不能直接输出裸字节数组,但可以通过自定义OutputFormat来绕过BytesWritable的字符串转换,直接写入原始字节流:
第一步:自定义RawOutputFormat
import org.apache.hadoop.fs.FSDataOutputStream; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.BytesWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.RecordWriter; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; public class RawOutputFormat extends FileOutputFormat<Text, BytesWritable> { @Override public RecordWriter<Text, BytesWritable> getRecordWriter(TaskAttemptContext job) throws IOException { Path outputDir = FileOutputFormat.getOutputPath(job); // 生成输出文件名,可根据TaskAttemptID动态生成避免冲突 Path outputFile = new Path(outputDir, "part-r-" + job.getTaskAttemptID().getId()); FileSystem fs = outputDir.getFileSystem(job.getConfiguration()); FSDataOutputStream outputStream = fs.create(outputFile); return new RecordWriter<>() { @Override public void write(Text key, BytesWritable value) throws IOException { // 直接写入字节数组的有效部分,跳过Writable的字符串转换 outputStream.write(value.getBytes(), 0, value.getLength()); // 如果需要分隔不同记录,可以添加自定义分隔符(比如换行符) // outputStream.write('\n'); } @Override public void close(TaskAttemptContext context) throws IOException { outputStream.close(); } }; } }
第二步:在MapReduce作业中配置该输出格式
Job job = Job.getInstance(conf, "RawByteOutputJob"); // ... 其他作业配置(设置Mapper、Reducer等) job.setOutputFormatClass(RawOutputFormat.class); // 设置输出键值类型(这里键用Text,值用BytesWritable) job.setOutputKeyClass(Text.class); job.setOutputValueClass(BytesWritable.class);
第三步:读取原始字节输出文件
此时输出文件是纯字节流,直接用普通文件输入流读取即可:
Configuration conf = new Configuration(); Path rawOutputPath = new Path("hdfs://your-output-path/part-r-00000"); FileSystem fs = rawOutputPath.getFileSystem(conf); try (FSDataInputStream inputStream = fs.open(rawOutputPath)) { byte[] buffer = new byte[4096]; int bytesRead; while ((bytesRead = inputStream.read(buffer)) != -1) { // 处理读取到的字节,buffer中前bytesRead个是有效数据 processRawBytes(buffer, 0, bytesRead); } } catch (IOException e) { e.printStackTrace(); }
内容的提问来源于stack exchange,提问作者Y0gesh Gupta
相关产品推荐
相关产品推荐

