序列化就是把内存中的对象,转换成字节序列(或其他数据传输协议)以便于存储到磁盘(持久化)和网络传输
反序列化就是将收到的字节序列(或其他数据传输协议)或者是磁盘的持久化数据,转换成内存中的对象
JDK自带的序列化只需要类实现了Serializable接口,就可以通过ObjectOutputStream类将对象变成byte[]字节数组。
所以,JDK自带的序列化在实际项目和框架中使用较少
实现bean对象序列化的步骤分为7步:
序列化中的顺序是无所谓的,但是序列化的顺序和反序列化的顺序必须是保持一致
如果,序列化的顺序是1,2,3;那么反序列化的顺序也得是1,2,3
序列化的顺序是2,1,3或者3,2,1都是可以的,反序列化的顺序与它保持一致即可
统计每一个手机号的总上行流量、总下行流量、总流量(总流量 = 总上行流量 + 总下行流量)
存在部分数据域名为空的情况
| id | 手机号码 | 网络ip | 域名 | 上行流量 | 下行流量 | 网络状态吗 |
|---|---|---|---|---|---|---|
| 7 | 13560436666 | 120.196.100.99 | www.abc.com | 1116 | 954 | 200 |
| 8 | 15910133277 | 192.168.100.54 | 315 | 296 | 200 |
| 手机号码 | 上行流量 | 下行流量 | 总流量 |
|---|---|---|---|
| 13560436666 | 1116 | 954 | 2070 |
package mapreduce.writable;
import org.apache.hadoop.io.Writable;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
public class FlowBean implements Writable {
private long upFlow;
private long downFlow;
private long sumFlow;
public FlowBean() {
}
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 void setSumFlow() {
this.sumFlow = this.upFlow + this.downFlow;
}
@Override
public String toString() {
return upFlow + "\t" + downFlow + "\t" + sumFlow;
}
@Override
public void write(DataOutput out) throws IOException {
out.writeLong(upFlow);
out.writeLong(downFlow);
out.writeLong(sumFlow);
}
@Override
public void readFields(DataInput in) throws IOException {
this.upFlow = in.readLong();
this.downFlow = in.readLong();
this.sumFlow = in.readLong();
}
}
package mapreduce.writable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import java.io.IOException;
public class FlowMapper extends Mapper<LongWritable, Text, Text, FlowBean> {
private Text outK = new Text();
private FlowBean outV = new FlowBean();
//1 13736230513 192.196.100.1 www.atguigu.com 2481 24681 200
@Override
protected void map(LongWritable key, Text value, Mapper<LongWritable, Text, Text, FlowBean>.Context context) throws IOException, InterruptedException {
String line = value.toString();
String[] split = line.split("\t");
String phone = split[1];
String upFlow = split[split.length - 3];
String downFlow = split[split.length - 2];
outK.set(phone);
outV.setUpFlow(Long.parseLong(upFlow));
outV.setDownFlow(Long.parseLong(downFlow));
outV.setSumFlow();
context.write(outK, outV);
}
}
package mapreduce.writable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import java.io.IOException;
public class FlowReducer extends Reducer<Text, FlowBean, Text, FlowBean> {
private long totalUpFlow;
private long totalDownFlow;
private FlowBean outV = new FlowBean();
@Override
protected void reduce(Text key, Iterable<FlowBean> values, Reducer<Text, FlowBean, Text, FlowBean>.Context context) throws IOException, InterruptedException {
totalUpFlow = 0;
totalDownFlow = 0;
for (FlowBean value : values) {
totalUpFlow += value.getUpFlow();
totalDownFlow += value.getDownFlow();
}
outV.setUpFlow(totalUpFlow);
outV.setDownFlow(totalDownFlow);
outV.setSumFlow();
context.write(key, outV);
}
}
package mapreduce.writable;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.io.IOException;
public class FlowDriver {
public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException {
Configuration configuration = new Configuration();
Job job = Job.getInstance(configuration);
job.setJarByClass(FlowDriver.class);
job.setMapperClass(FlowMapper.class);
job.setReducerClass(FlowReducer.class);
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(FlowBean.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(FlowBean.class);
FileInputFormat.setInputPaths(job, new Path("D:\\input\\writable"));
FileOutputFormat.setOutputPath(job, new Path("D:\\output\\writable"));
boolean res = job.waitForCompletion(true);
System.exit(res ? 0 : 1);
}
}
下图是输出结果

下图是输入数据
两个蓝色框中的手机号是一样的,检查一下代码结果的正确性

上传到集群的代码,在获取输入路径和输出路径时,应该从命令中读取,所以Driver类代码要做一点小修改
//设置输入和输出路径
FileInputFormat.setInputPaths(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
[stu@hadoop102 hadoop-3.1.3]$ hadoop fs -mkdir /input
[stu@hadoop102 hadoop-3.1.3]$ hadoop fs -put phone_data.txt /input
需要在pom.xml中添加打包插件依赖
<build>
<plugins>
<plugin>
<artifactId>maven-compiler-pluginartifactId>
<version>3.6.1version>
<configuration>
<source>1.8source>
<target>1.8target>
configuration>
plugin>
<plugin>
<artifactId>maven-assembly-pluginartifactId>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependenciesdescriptorRef>
descriptorRefs>
configuration>
<executions>
<execution>
<id>make-assemblyid>
<phase>packagephase>
<goals>
<goal>singlegoal>
goals>
execution>
executions>
plugin>
plugins>
build>
在target目录下,生成了两个jar包,选择没有jar包依赖的


上传到/opt/module/hadoop-3.1.3下

[stu@hadoop102 hadoop-3.1.3]$ hadoop jar writable.jar mapreduce.writable.FlowDriver /input/phone_data.txt /output

查看part-r-00000的文件内容,和在本地运行时的结果一致
