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

Hadoop学习---第四篇Mapreducer里的Partitioner,hadooppartitioner

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

本文主要包含Hadoop学习---第四篇Mapreducer里的Partitioner,hadooppartitioner等服务器相关知识,网友希望可以进行参考

Hadoop学习---第四篇Mapreducer里的Partitioner,hadooppartitioner


Partitioner就是对map输出的key进行分组,不同的组可以指定不同的reduce task处理;

Partition功能由partitioner的实现子类来实现

每写一段代码都会加深理解,程序里记录了自己的理解

FlowBean类源码:

package cn.zxl.flowcountpartitioner;

import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;

import org.apache.hadoop.io.WritableComparable;

public class FlowBean implements WritableComparable<FlowBean>{
	private long upflow;//上行流量
	private long downflow;//下行流量
	private long sumflow;//总流量
	public long getUpflow() {
		return upflow;
	}
	public void setUpflow(long upflow) {
		this.upflow = upflow;
	}
	public long getDownflow() {
		return downflow;
	}
	public void setDownflow(long downflow) {
		this.downflow = downflow;
	}
	public long getSumflow() {
		return sumflow;
	}
	public void setSumflow(long sumflow) {
		this.sumflow = sumflow;
	}
	public FlowBean() {
	}
	public FlowBean(long upflow, long downflow) {
		super();
		this.upflow = upflow;
		this.downflow = downflow;
		this.sumflow = upflow+downflow;
	}
	@Override
	public void readFields(DataInput in) throws IOException {
		upflow=in.readLong();
		downflow=in.readLong();
		sumflow=in.readLong();
	}
	@Override
	public void write(DataOutput out) throws IOException {
		out.writeLong(upflow);
		out.writeLong(downflow);
		out.writeLong(sumflow);
		
	}
	@Override
	public int compareTo(FlowBean bean) {
		return sumflow>bean.getSumflow()?-1:1;
	}
	
	@Override
	public String toString() {
		return upflow+"\t"+downflow+"\t"+sumflow;
	}
}
ProvicePartition类源码:

package cn.zxl.flowcountpartitioner;

import java.util.HashMap;

import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Partitioner;

public class ProvicePartition extends Partitioner<Text, FlowBean>{
	//根据手机号前三位划分分组
	//Partitioner就是对key进行分组
	private static HashMap<String, Integer> pmap = new HashMap<String, Integer>();
	static{
		pmap.put("136", 0);
		pmap.put("137", 1);
		pmap.put("138", 2);
		pmap.put("139", 3);
	}
	@Override
	public int getPartition(Text key, FlowBean bean, int numPartitions) {
		String prex=key.toString().substring(0,3);
		Integer partNum=pmap.get(prex);//根据key截取的前三位做key和map的值是否匹配
		return partNum==null?4:partNum;
	}
}

FlowCount类源码:

package cn.zxl.flowcountpartitioner;

import java.io.IOException;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;

public class FlowCount {
	static class FlowCountMapper extends Mapper<LongWritable, Text, Text, FlowBean>{
		@Override
		protected void map(LongWritable key, Text value,Context context)
				throws IOException, InterruptedException {
			String line=value.toString();
			String[] phoneinfo=line.split("\t");
			String phoneN=phoneinfo[0];
			String upflow=phoneinfo[phoneinfo.length-3];
			String downflow=phoneinfo[phoneinfo.length-2];
			FlowBean fb=new FlowBean(Long.parseLong(upflow),Long.parseLong(downflow));
			context.write(new Text(phoneN),fb);
		}
	}
	//reducer里的值是<key,list(value)>,也就是相同的键里对应一个集合
	//reducer是根据key排序的,而不是value,要根据什么排序,那就得已什么作为key输出
	static class FlowCountReducer extends Reducer<Text, FlowBean, Text, FlowBean>{
		@Override
		protected void reduce(Text key, Iterable<FlowBean> values,Context context)
				throws IOException, InterruptedException {
			long upflow_sum=0;
			long downflow_sum=0;
			for(FlowBean bean:values){
				upflow_sum+=bean.getUpflow();
				downflow_sum+=bean.getDownflow();
			}
			FlowBean fb=new FlowBean(upflow_sum,downflow_sum);
			context.write(new Text(key), fb);
		}
	}
	
	public static void main(String[] args) throws Exception {
		Configuration conf=new Configuration();// cn.zxl.flowcountpartitioner.FlowCount
		
		Job job=Job.getInstance(conf);
		
		job.setJarByClass(FlowCount.class);
		
		job.setMapperClass(FlowCountMapper.class);
		job.setReducerClass(FlowCountReducer.class);
		
		job.setOutputKeyClass(Text.class);
		job.setOutputValueClass(FlowBean.class);
		//job指定自定义的Partitioner组件
		//job.setPartitionerClass(ProvicePartition.class);
		/*job中指定reducertask的数量,说明:这里的reducertask数量可以指定为1个,如果是1个reducertask,
		*那么所有的分区数据都输入到一个文件里,如果指定个数小于分区个数(这里是5个),那么程序会报错,
		*因为不知道对应的一个分区数据放置到哪里,如果指定个数超过分区个数,那么后面产生的文件是空的
		*/
		//job.setNumReduceTasks(5);
		
		FileInputFormat.setInputPaths(job, new Path(args[0]));
		
		Path output = new Path(args[1]);
		FileSystem fs = output.getFileSystem(conf);
		//看输出是否存在,存在就删除,特别说明:安全起见正式的线上建议最好不要做这个判断,如果这样做,会把以前产生的数据删除
		//补充:正式生产环境最好指定删除多久后正式删除数据,以便错删时可以恢复数据
		/*
		 * 添加在hdfs-site的配置文件里
		 * <property>
			<name>fs.trash.interval</name>
			<value>60</value><!-- 回收站过期机制检查频率(分钟) -->
			</property>
			
			<property>
			<name>fs.trash.checkpoint.interval</name>
			<value>20</value><!-- 回收站中文件过期的时间限制(分钟) -->
			</property>
		 */
		if(fs.exists(output)){
			fs.delete(output, true);
		}
		FileOutputFormat.setOutputPath(job, new Path(args[1]));
		
		job.waitForCompletion(true);
	}
}

测试数据:

1363157985066 1372623050300-FD-07-A4-72-B8:CMCC 120.196.100.82i02.c.aliimg.com 2427 248124681 200
1363157995052 138265441015C-0E-8B-C7-F1-E0:CMCC 120.197.40.44 0 264 0 200
1363157991076 1392643565620-10-7A-28-CC-0A:CMCC 120.196.100.992 4 132 1512 200
1363154400022 139262511065C-0E-8B-8B-B1-50:CMCC 120.197.40.44 0 240 0 200
1363157993044 1821157596194-71-AC-CD-E6-18:CMCC-EASY 120.196.100.99iface.qiyi.com 视频网站15 12 1527 2106 200
1363157995074 841384135C-0E-8B-8C-E8-20:7DaysInn 120.197.40.4122.72.52.12 2016 41161432 200
1363157993055 13560439658C4-17-FE-BA-DE-D9:CMCC 120.196.100.9918 15 1116 954 200
1363157995033 159201332575C-0E-8B-C7-BA-20:CMCC 120.197.40.4sug.so.360.cn 信息安全20 20 3156 2936 200
1363157983019 1371919941968-A1-B7-03-07-B1:CMCC-EASY<

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

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

  • Hadoop学习---第四篇Mapreducer里的Partitioner,hadooppartitioner

相关文章

  • Xen安全架构sHype/ACM和XSM/Flask的相关网络资源,xenshype
  • MapReduce程序之实现单表关联,mapreduce单表关联
  • &lt;转&gt;Openstack ceilometer 宿主机监控模块扩展,openstackceilometer
  • 微软小冰、小娜不久相会在中国,微软相会在中国
  • pagerank算法的MapReduce实现,pagerankmapreduce
  • @RequestMapping的params参数 @RequestMapping(params = &quot;method=save&quot;),requestmapping
  • 在linux下制作libxxx.so 动态库,linuxlibxxx.so
  • 常见类型网站的SEO应该怎么操作?类型不同,操作自然不同,因为分门别类嘛,seo分门别类
  • Hive 外部表 分区表,hive外部表分区表
  • elasticsearch JAVA客户端操作---搜索的过滤、分组高亮,elasticsearchjava

文章分类

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

最近更新的内容

    • Greenplum+Hadoop学习笔记-15-管理数据,hadoop-15-
    • Scala类型详解,scala详解
    • 程序员必备的代码审查(Code Review)清单,codereview
    • Hive快速入门,hive入门
    • 云主机跟VPS哪个比较好?哪个稳定安全?,主机vps
    • Xen安全架构sHype/ACM和XSM/Flask的相关网络资源,xenshype
    • hive的变量传递设置,hive变量传递设置
    • Greenplum数据库升级实务(上),greenplum实务
    • Deploy FAILURE:An app Was not successfully detected by any available buildpack,deploybuildpack
    • MapReduce编程之实现多表关联,mapreduce编程关联

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

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