• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > MapReduce之RecordReader组件源码解析及实例,mapreduce实例

MapReduce之RecordReader组件源码解析及实例,mapreduce实例

作者:网友 字体:[增加 减小] 来源:互联网

本文主要包含node mapreduce,mapreduce是什么,mapreduce,mapreduce原理,mapreduce编程实例等服务器相关知识,网友希望可以进行参考

MapReduce之RecordReader组件源码解析及实例,mapreduce实例


简述

无论我们以怎样的方式从分片中读取一条记录,每读取一条记录都会调用RecordReader类;
系统默认的RecordReader是LineRecordReader,TextInputFormat;
LineRecordReader是用每行的偏移量作为map的key,每行的内容作为map的value;
而SequenceFileInputFormat的RecordReader是SequenceFileRecordReader;
应用场景:自定义读取每一条记录的方式;自定义读入key的类型,如希望读取的key是文件的路径或名字而不是该行在文件中的偏移量。

TextInputFormat源码如下:

package org.apache.hadoop.mapreduce.lib.input;
/** An {@link InputFormat} for plain text files.  Files are broken into lines.
 * Either linefeed or carriage-return are used to signal end of line.  Keys are
 * the position in the file, and values are the line of text.. */
public class TextInputFormat extends FileInputFormat<LongWritable, Text> {

  @Override
  public RecordReader<LongWritable, Text> 
    createRecordReader(InputSplit split,
                       TaskAttemptContext context) {
    String delimiter = context.getConfiguration().get(
        "textinputformat.record.delimiter");
    byte[] recordDelimiterBytes = null;
    if (null != delimiter)
      recordDelimiterBytes = delimiter.getBytes(Charsets.UTF_8);
    return new LineRecordReader(recordDelimiterBytes);
  }

  @Override
  protected boolean isSplitable(JobContext context, Path file) {
    final CompressionCodec codec =
      new CompressionCodecFactory(context.getConfiguration()).getCodec(file);
    if (null == codec) {
      return true;
    }
    return codec instanceof SplittableCompressionCodec;
  }
}

textinputformat.record.delimiter指的是读取一行的数据的终止符号,即遇到textinputformat.record.delimiter所包含的字符时,该一行的读取结束。
可以通过Configuration的set()方法来设置自定义的终止符,如果没有设置 textinputformat.record.delimiter,那么Hadoop就采用以CR,LF或者CRLF作为终止符,这一点可以查看LineReader的readDefaultLine方法 。

LineRecordReader源码如下:

package org.apache.hadoop.mapreduce.lib.input;

/**
 * Treats keys as offset in file and value as line. 
 */
public class LineRecordReader extends RecordReader<LongWritable, Text> {
     public void initialize(InputSplit genericSplit,
                         TaskAttemptContext context) throws IOException {
        ......
        start = split.getStart();
        end = start + split.getLength();
        final Path file = split.getPath();

        // open the file and seek to the start of the split
        final FileSystem fs = file.getFileSystem(job);
        fileIn = fs.open(file);
        ......
        // If this is not the first split, we always throw away first record
        // because we always (except the last split) read one extra line in
        // next() method.
        if (start != 0) {
            start += in.readLine(new Text(), 0, maxBytesToConsume(start));
        }
        this.pos = start;
        ......
    }
    ......
}

自定义RecordReader

1、继承抽象类RecordReader,实现RecordReader的一个实例;
2、实现自定义InputFormat类,重写InputFormat中createRecordReader()方法,返回值是自定义的RecordReader实例;
3、配置job.setInputFormatClass()设置自定义的InputFormat实例;

实例

数据:

10
20
30
40
50
60
70
……

要求:读取整个文件,分别计算奇数行与偶数行数据之和
奇数行之和:10+30+50+70=160
偶数行之和:20+40+60=120

package Recordreader;

import java.io.IOException;
import java.net.URI;
import java.net.URISyntaxException;
import java.util.List;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataInputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.InputSplit;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.JobContext;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Partitioner;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.FileSplit;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.LineReader;

public class MyRecordReader {

    private final static String INPUT_PATH = "hdfs://liguodong:8020/inputsum";
    private final static String OUTPUT_PATH = "hdfs://liguodong:8020/outputsum";    

    public static class DefRecordReader extends 


 
分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • MapReduce之RecordReader组件源码解析及实例,mapreduce实例
  • Java MapReduce详解--(3),mapreduce详解
  • MapReduce对输入多文件的处理,mapreduce输入处理
  • MapReduce处理表的自连接,mapreduce处理表
  • MapReduce的两表join操作优化,mapreduce表join
  • MapReduce的两表join一般操作,mapreduce表join

相关文章

  • Hadoop集群优化,hadoop集群
  • Cloud Foundry service broker开发部署实例解析(下),foundrybroker
  • HIVE以及OOZIE添加第三方JAR包的方法,hiveoozie
  • kafka 效率优化,kafka优化
  • Storm On Yarn部署,stormonyarn部署
  • HDFS小文件的合并优化
  • Libvirt中windows虚拟机的动态内存管理,libvirt虚拟机
  • A Note on Distributed Computing,anoteondialectic
  • 机器学习--Logistic回归算法案例,--logistic算法
  • R.net简介(原创翻译),r.net简介原创翻译

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • (Kilo)Devstack Kilo版本localrc推荐,devstacklocalrc
    • Greenplum+Hadoop学习笔记-14-定义数据库对象之创建与管理表空间,hadoop-14-
    • spark streaming 调试技巧,sparkstreaming
    • spark core源码分析15 Shuffle详解-写流程,sparkshuffle
    • hadoop集群扩展,hadoop集群
    • hadoop-common源码分析之-WritableUtils,hadoopcommon源码
    • Apache Pig的前世今生,apachepig前世今生
    • Spark下的PageRank实现,SparkPageRank实现
    • 在Java程序中调用Salesforce REST API,salesforcerest
    • hadoop-2.6.0 Unhealthy Nodes 问题,hadoop2.6.0安装

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有