Hadoop_25_MapReduce实现日志清洗程序

1、需求:

对web访问日志中的各字段识别切分,去除日志中不合法的记录,根据KPI统计需求,生成各类访问请求过滤数据

2、实现代码:

a) 定义一个bean,用来记录日志数据中的各数据字段

1package cn.bigdta.hdfs.weblog; 2 3public class WebLogBean { 4 5 private String remote_addr;// 记录客户端的ip地址 6 private String remote_user;// 记录客户端用户名称,忽略属性"-" 7 private String time_local;// 记录访问时间与时区 8 private String request;// 记录请求的url与http协议 9 private String status;// 记录请求状态;成功是200 10 private String body_bytes_sent;// 记录发送给客户端文件主体内容大小 11 private String http_referer;// 用来记录从那个页面链接访问过来的 12 private String http_user_agent;// 记录客户浏览器的相关信息 13 14 private boolean valid = true;// 判断数据是否合法 15 16 17 18 public String getRemote_addr() { 19 return remote_addr; 20 } 21 22 public void setRemote_addr(String remote_addr) { 23 this.remote_addr = remote_addr; 24 } 25 26 public String getRemote_user() { 27 return remote_user; 28 } 29 30 public void setRemote_user(String remote_user) { 31 this.remote_user = remote_user; 32 } 33 34 public String getTime_local() { 35 return time_local; 36 } 37 38 public void setTime_local(String time_local) { 39 this.time_local = time_local; 40 } 41 42 public String getRequest() { 43 return request; 44 } 45 46 public void setRequest(String request) { 47 this.request = request; 48 } 49 50 public String getStatus() { 51 return status; 52 } 53 54 public void setStatus(String status) { 55 this.status = status; 56 } 57 58 public String getBody_bytes_sent() { 59 return body_bytes_sent; 60 } 61 62 public void setBody_bytes_sent(String body_bytes_sent) { 63 this.body_bytes_sent = body_bytes_sent; 64 } 65 66 public String getHttp_referer() { 67 return http_referer; 68 } 69 70 public void setHttp_referer(String http_referer) { 71 this.http_referer = http_referer; 72 } 73 74 public String getHttp_user_agent() { 75 return http_user_agent; 76 } 77 78 public void setHttp_user_agent(String http_user_agent) { 79 this.http_user_agent = http_user_agent; 80 } 81 82 public boolean isValid() { 83 return valid; 84 } 85 86 public void setValid(boolean valid) { 87 this.valid = valid; 88 } 89 90 91 @Override 92 public String toString() { 93 StringBuilder sb = new StringBuilder(); 94 sb.append(this.valid); 95 sb.append("\001").append(this.remote_addr); 96 sb.append("\001").append(this.remote_user); 97 sb.append("\001").append(this.time_local); 98 sb.append("\001").append(this.request); 99 sb.append("\001").append(this.status); 100 sb.append("\001").append(this.body_bytes_sent); 101 sb.append("\001").append(this.http_referer); 102 sb.append("\001").append(this.http_user_agent); 103 return sb.toString(); 104 } 105}

View Code

 b)定义一个parser用来解析过滤web访问日志原始记录

1package cn.bigdta.hdfs.weblog; 2 3import java.text.ParseException; 4import java.text.SimpleDateFormat; 5import java.util.Date; 6import java.util.Locale; 7 8public class WebLogParser { 9 10 static SimpleDateFormat sd1 = new SimpleDateFormat("dd/MMM/yyyy:HH:mm:ss", Locale.US); 11 12 static SimpleDateFormat sd2 = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); 13 14 public static WebLogBean parser(String line) { 15 WebLogBean webLogBean = new WebLogBean(); 16 String[] arr = line.split(" "); 17 if (arr.length > 11) { 18 webLogBean.setRemote_addr(arr[0]); 19 webLogBean.setRemote_user(arr[1]); 20 webLogBean.setTime_local(parseTime(arr[3].substring(1))); 21 webLogBean.setRequest(arr[6]); 22 webLogBean.setStatus(arr[8]); 23 webLogBean.setBody_bytes_sent(arr[9]); 24 webLogBean.setHttp_referer(arr[10]); 25 26 if (arr.length > 12) { 27 webLogBean.setHttp_user_agent(arr[11] + " " + arr[12]); 28 } else { 29 webLogBean.setHttp_user_agent(arr[11]); 30 } 31 if (Integer.parseInt(webLogBean.getStatus()) >= 400) {// 大于400,HTTP错误 32 webLogBean.setValid(false); 33 } 34 } else { 35 webLogBean.setValid(false); 36 } 37 return webLogBean; 38 } 39 40 public static String parseTime(String dt) { 41 42 String timeString = ""; 43 try { 44 Date parse = sd1.parse(dt); 45 timeString = sd2.format(parse); 46 47 } catch (ParseException e) { 48 e.printStackTrace(); 49 } 50 return timeString; 51 } 52}

View Code

 c) mapreduce程序

1package cn.bigdta.hdfs.weblog; 2import java.io.IOException; 3import org.apache.hadoop.conf.Configuration; 4import org.apache.hadoop.fs.Path; 5import org.apache.hadoop.io.LongWritable; 6import org.apache.hadoop.io.NullWritable; 7import org.apache.hadoop.io.Text; 8import org.apache.hadoop.mapreduce.Job; 9import org.apache.hadoop.mapreduce.Mapper; 10import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; 11import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; 12 13public class WeblogPreProcess { 14 15 static class WeblogPreProcessMapper extends Mapper<LongWritable, Text, Text, NullWritable> { 16 Text k = new Text(); 17 NullWritable v = NullWritable.get(); 18 19 @Override 20 protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { 21 22 String line = value.toString(); 23 WebLogBean webLogBean = WebLogParser.parser(line); 24 //可以插入一个静态资源过滤(.....) 25 /*WebLogParser.filterStaticResource(webLogBean);*/ 26 if (!webLogBean.isValid()) 27 return; 28 k.set(webLogBean.toString()); 29 context.write(k, v); 30 } 31 } 32 33 public static void main(String[] args) throws Exception { 34 35 Configuration conf = new Configuration(); 36 Job job = Job.getInstance(conf); 37 38 job.setJarByClass(WeblogPreProcess.class); 39 40 job.setMapperClass(WeblogPreProcessMapper.class); 41 42 job.setOutputKeyClass(Text.class); 43 job.setOutputValueClass(NullWritable.class); 44 45 FileInputFormat.setInputPaths(job, new Path("F:/weblog")); 46 FileOutputFormat.setOutputPath(job, new Path("F:/weblogOut")); 47 48 job.waitForCompletion(true); 49 50 } 51}

 日志文件下载:https://pan.baidu.com/s/17oOaA_S5RRDKjFhCqMm40A

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

Hadoop_25_MapReduce实现日志清洗程序 - HelloWorld