在上一篇中介紹了MapReduce進行單詞計數的案例,這一章介紹下怎么查看MapReduce的任務日志,
如果想要查看mapreduce任務執行程序產生的日志資訊怎么辦呢?
是不是在提交任務的時候直接在這個控制臺上就能看到了?先不要著急,我們先在代碼中增加一些日志資訊,在實際作業中做除錯的時候這個也是很有必要的
一、syout日志輸出
1、mapper類修改
在自定義mapper類的map函式中增加一個輸出,將k1,v1的值列印出來
添加內容如下:

mapper類修改后代碼如下:
public static class MyMapper extends Mapper<LongWritable,Text,Text,LongWritable>{
/**
* 需要實作map函式
* 這個map函式就是可以接收k1,v1, 產生k2,v2
* @param k1
* @param v1
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void map(LongWritable k1, Text v1, Context context) throws IOException, InterruptedException {
//輸出k1,v1的值
System.out.println("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
//k1代表的是每一行的行首偏移量,v1代表的是每一行的內容
//對獲取到的每一行資料進行切割,把單詞切割出來
String[] words = v1.toString().split(" ");
//迭代切割出來的單詞資料
for(String word:words){
//把迭代出來的單詞封裝成<k2,v2>的形式
Text k2 = new Text(word);
LongWritable v2 = new LongWritable(1L);
System.out.println("k2:"+word+"...v2:1");
//把<k2,v2>寫出去
context.write(k2,v2);
}
}
}
2、reducer類修改
在自定義reducer類中的reduce方法中增加一個輸出,將k2,v2和k3,v3的值列印出來

修改后reducer代碼如下:
public static class MyReducer extends Reducer<Text,LongWritable,Text,LongWritable>{
/**
* 針對v2s的資料進行累加求和,并且最終把資料轉化為k3,v3寫出去
* @param k2
* @param v2s
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void reduce(Text k2, Iterable<LongWritable> v2s, Context context) throws IOException, InterruptedException {
//創建一個sum變數,保持v2s的和
long sum = 0L;
for (LongWritable v2:v2s){
//輸出k2,v2的值
System.out.println("<k2,v2>=<"+k2.toString()+","+v2.get()+">");
sum += v2.get();
}
//組裝k3,v3
Text k3 = k2;
LongWritable v3 = new LongWritable(sum);
//輸出k3,v3的值
System.out.println("<k3,v3>=<"+k3.toString()+","+v3.get()+">");
//把結果寫出去
context.write(k3,v3);
}
}
3、打包
把之前的jar包改個名字備份一下
mv db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar_01

重新在windows機器上打jar包
mvn clean package -DskipTests

我們發現報錯了,因為上傳到服務器占用了這個目錄!!!
這里的解決辦法是隨便切到一個目錄,上傳一個東西這樣就不會占用了!!關閉MobaXterm重新打開MobaXterm是沒有用的,還是會記錄上一次的目錄!!

我這里重新在/data/soft目錄上傳了一個hello.txt檔案,

這樣的話,本地專案所在的目錄就不會被占用了
E:\project\db_hadoop\target
然后,重新執行如下命令
mvn clean package -DskipTests

這樣就能打包成功了,
把新的jar包上傳到bigdata01機器的/data/soft/jar目錄中

下面那個是我們備份的,不用管它,
4、提交任務
重新向集群提交任務,注意,針對輸出目錄,要么換一個新的不存在的目錄,要么把之前的out目錄刪掉可以繼續使用out目錄,在這里我換了一個新的輸出目錄 out1
hadoop jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar com.imooc.mr.WordCountJob /test/hello.txt /out1
等待任務執行結束,我們發現在控制臺上是看不到任務中的日志資訊的,為什么呢?因為我們在這相當于是通過一個客戶端把任務提交到集群里面去執行了,所以日志是存在在集群里面的,想要查看需要需要到一個特殊的地方查看這些日志資訊,

先進入到yarn的web界面,訪問8088埠,點擊對應任務的history鏈接
http://bigdata01:8088/
注意:如果想使用主機名在瀏覽器中訪問8088界面,不能開啟翻墻工具,否則就算在windows的hosts檔案中配置了虛擬機的ip和主機名的映射關系,在瀏覽器中也無法使用主機名訪問,因為此時請求會被翻墻工具所攔截,但是翻墻工具中無法識別這個主機名,

注意了,在這里我們發現這個鏈接是打不來的,如下

原因有2個:
第1個原因:
沒有在本地C:\Windows\System32\drivers\etc\hosts進行配置主機名和ip的映射關系,
192.168.18.100 bigdata01

但我們發現我們已經進行配置了,

第2個原因:
就是這里必須要啟動historyserver行程才可以,并且還要開啟日志聚合功能,才能在web界面上直接查看任務對應的日志資訊,因為默認情況下任務的日志是散落在nodemanager節點上的,想要查看需要找到對應的nodemanager節點上去查看,這樣就很不方便,通過日志聚合功能我們可以把之前本來散落在nodemanager節點上的日志統一收集到hdfs上的指定目錄中,這樣就可以在yarn的web界面中直接查看了,
那我們就來開啟日志聚合功能,開啟日志聚合功能需要修改yarn-site.xml的配置,增加
yarn.log-aggregation-enable和yarn.log.server.url這兩個引數
<property>
<name>yarn.log-aggregation-enable</name>
<value>true</value>
</property>
<property>
<name>yarn.log.server.url</name>
<value>http://bigdata01:19888/jobhistory/logs/</value>
</property>
然后我們進行修改
cd /data/soft/hadoop-3.2.0/etc/hadoop
vi yarn-site.xml
添加后如下:

注意,修改完成后重新啟動集群,
stop-all.sh
start-all.sh

啟動historyserver行程,需要在集群的所有節點上都啟動這個行程,因為我這里是偽分布部署的,只有一個節點,只需要啟動一個,
mapred --daemon start historyserver

5、重新提交任務
重啟集群后之前的任務就沒了,這里我重新修改了輸出目錄為/out2
cd /data/soft/jar/
hadoop jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar com.imooc.mr.WordCountJob /test/hello.txt /out2
等待任務執行完成


此時再進入yarn的8088界面,點擊任務對應的history鏈接就可以打開了,

此時,點擊對應map和reduce后面的鏈接就可以點進去查看日志資訊了,點擊map后面的數字1,可以進入如下界面


點擊這個界面中的logs文字鏈接,可以查看詳細的日志資訊,

最終可以在界面中看到很多日志資訊,我們剛才使用sout輸出的日志資訊需要到Log Type: stdout這里來查看,在這里可以看到,k1和v1的值

想要查看reduce輸出的日志資訊需要到reduce里面查看,操作流程是一樣的,可以看到k2,v2和k3,v3的值,




咱們剛才的輸出是使用syout輸出的,這個其實是不正規的,標準的日志寫法是需要使用logger進行輸出的,
二、logger日志輸出
1、修改mapper類
修改如下部分:

修改后代碼如下:
public static class MyMapper extends Mapper<LongWritable,Text,Text,LongWritable>{
Logger logger = LoggerFactory.getLogger(MyMapper.class);
/**
* 需要實作map函式
* 這個map函式就是可以接收k1,v1, 產生k2,v2
* @param k1
* @param v1
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void map(LongWritable k1, Text v1, Context context) throws IOException, InterruptedException {
//輸出k1,v1的值
//System.out.println("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
logger.info("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
//k1代表的是每一行的行首偏移量,v1代表的是每一行的內容
//對獲取到的每一行資料進行切割,把單詞切割出來
String[] words = v1.toString().split(" ");
//迭代切割出來的單詞資料
for(String word:words){
//把迭代出來的單詞封裝成<k2,v2>的形式
Text k2 = new Text(word);
LongWritable v2 = new LongWritable(1L);
//System.out.println("k2:"+word+"...v2:1");
//把<k2,v2>寫出去
context.write(k2,v2);
}
}
}
2、修改reducer類
修改如下部分:

修改后代碼如下:
public static class MyReducer extends Reducer<Text,LongWritable,Text,LongWritable>{
Logger logger = LoggerFactory.getLogger(MyReducer.class);
/**
* 針對v2s的資料進行累加求和,并且最終把資料轉化為k3,v3寫出去
* @param k2
* @param v2s
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void reduce(Text k2, Iterable<LongWritable> v2s, Context context) throws IOException, InterruptedException {
//創建一個sum變數,保持v2s的和
long sum = 0L;
for (LongWritable v2:v2s){
//輸出k2,v2的值
//System.out.println("<k2,v2>=<"+k2.toString()+","+v2.get()+">");
logger.info("<k2,v2>=<"+k2.toString()+","+v2.get()+">");
sum += v2.get();
}
//組裝k3,v3
Text k3 = k2;
LongWritable v3 = new LongWritable(sum);
//輸出k3,v3的值
//System.out.println("<k3,v3>=<"+k3.toString()+","+v3.get()+">");
logger.info("<k3,v3>=<"+k3.toString()+","+v3.get()+">");
//把結果寫出去
context.write(k3,v3);
}
}
修改后整個代碼匯總如下:
package com.imooc.mr;
import org.apache.hadoop.conf.Configuration;
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;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
/**
* 需求:讀取hdfs上的hello.txt檔案,計算檔案中每個單詞出現的總次數
* hello.txt檔案內容如下:
* hello you
* hello me
* 最終需要的結果形式如下:
* hello 2
* me 1
* you 1
*/
public class WordCountJob {
/**
* 組裝job=map+reduce
* @param args
*/
public static void main(String[] args) {
try {
if(args.length != 2){
System.exit(100);
}
//job需要的配置引數
Configuration conf = new Configuration();
//創建一個job
Job job = Job.getInstance(conf);
// 注意:這一行必須設定,否則在集群中執行的是找不到WordCountJob這個類
job.setJarByClass(WordCountJob.class);
//指定輸入路徑(可以是檔案,也可以是目錄)
FileInputFormat.setInputPaths(job,new Path(args[0]));
//指定輸出路徑(只能指定一個不存在的目錄)
FileOutputFormat.setOutputPath(job,new Path(args[1]));
//指定map相關的代碼
job.setMapperClass(MyMapper.class);
//指定k2的型別
job.setMapOutputKeyClass(Text.class);
//指定v2的型別
job.setMapOutputValueClass(LongWritable.class);
//指定reduce相關的代碼
job.setReducerClass(MyReducer.class);
//指定k3的型別
job.setOutputKeyClass(Text.class);
//指定v3的型別
job.setOutputValueClass(LongWritable.class);
//提交job
job.waitForCompletion(true);
}catch (Exception e){
e.printStackTrace();
}
}
public static class MyMapper extends Mapper<LongWritable,Text,Text,LongWritable>{
Logger logger = LoggerFactory.getLogger(MyMapper.class);
/**
* 需要實作map函式
* 這個map函式就是可以接收k1,v1, 產生k2,v2
* @param k1
* @param v1
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void map(LongWritable k1, Text v1, Context context) throws IOException, InterruptedException {
//輸出k1,v1的值
//System.out.println("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
logger.info("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
//k1代表的是每一行的行首偏移量,v1代表的是每一行的內容
//對獲取到的每一行資料進行切割,把單詞切割出來
String[] words = v1.toString().split(" ");
//迭代切割出來的單詞資料
for(String word:words){
//把迭代出來的單詞封裝成<k2,v2>的形式
Text k2 = new Text(word);
LongWritable v2 = new LongWritable(1L);
//System.out.println("k2:"+word+"...v2:1");
//把<k2,v2>寫出去
context.write(k2,v2);
}
}
}
public static class MyReducer extends Reducer<Text,LongWritable,Text,LongWritable>{
Logger logger = LoggerFactory.getLogger(MyReducer.class);
/**
* 針對v2s的資料進行累加求和,并且最終把資料轉化為k3,v3寫出去
* @param k2
* @param v2s
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void reduce(Text k2, Iterable<LongWritable> v2s, Context context) throws IOException, InterruptedException {
//創建一個sum變量,保持v2s的和
long sum = 0L;
for (LongWritable v2:v2s){
//輸出k2,v2的值
//System.out.println("<k2,v2>=<"+k2.toString()+","+v2.get()+">");
logger.info("<k2,v2>=<"+k2.toString()+","+v2.get()+">");
sum += v2.get();
}
//組裝k3,v3
Text k3 = k2;
LongWritable v3 = new LongWritable(sum);
//輸出k3,v3的值
//System.out.println("<k3,v3>=<"+k3.toString()+","+v3.get()+">");
logger.info("<k3,v3>=<"+k3.toString()+","+v3.get()+">");
//把結果寫出去
context.write(k3,v3);
}
}
}
重新編譯打包上傳,重新提交最新的jar包,這個時候再查看日志就需要到Log Type: syslog中查看日志了,
3、打包
打包失敗,一般是目錄被占用,按照前面說的方法去處理,這里就不演示了,
mvn clean package -DskipTests

4、上傳
先備份下之前jar包
cd /data/soft/jar
mv db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar_02
然后上傳jar包

5、提交任務
hadoop jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar com.imooc.mr.WordCountJob /test/hello.txt /out3

6、查看日志
這個時候再查看日志就需要到Log Type: syslog中查看日志了,
點擊History

(1)map的日志:

點擊logs

找到Log Type: syslog


(2)reduce的日志:

點擊logs

找到Log Type: syslog


這是作業中比較常用的查看日志的方式
三、其他方式查看日志
但是還有一種使用命令查看的方式,這種方式面試的時候一般喜歡問,如下:
yarn logs -applicationId application_1646133121455_0003

注意:后面指定的是任務id,任務id可以到yarn的web界面上查看,

執行這個命令可以看到很多的日志資訊,我們通過grep篩選一下日志,
yarn logs -applicationId application_1646133121455_0003 | grep k1,v1

yarn logs -applicationId application_1646133121455_0003 | grep k2,v2
yarn logs -applicationId application_1646133121455_0003 | grep k3,v3

這種方式也需要大家能夠記住并且掌握住,首先是面試的時候可能會問到,還有就是針對某一些艱難的場景下,無法使用yarn的web界面查看日志,就需要使用yarn logs命令了,
補充2個內容:
四、停止Hadoop集群中的任務
如果一個mapreduce任務處理的資料量比較大的話,這個任務會執行很長時間,可能幾十分鐘或者幾個小時都有可能,假設一個場景,任務執行了一半了我們發現我們的代碼寫的有問題,需要修改代碼重新提交執行,這個時候之前的任務就沒有必要再執行了,沒有任何意義了,最終的結果肯定是錯誤的,所以我們就想把它停掉,要不然會額外浪費集群的資源,如何停止呢?
我在提交任務的視窗中按ctrl+c是不是就可以停止?
注意了,不是這樣的,我們前面說過,這個任務是提交到集群執行的,你在提交任務的視窗中執行ctrl+c對已經提交到集群中的任務是沒有任何影響的,
我們可以驗證一下,執行ctrl+c之后你再到yarn的8088界面查看,會發現任務依然存在,
所以需要使用hadoop集群的命令去停止正在運行的任務
使用yarn application -kill命令,后面指定任務id即可
yarn application -kill application_1587713567839_0003
五、MapReduce程式擴展
咱們前面說過MapReduce任務是由map階段和reduce階段組成的
但是我們也說過,reduce階段不是必須的,那也就意味著MapReduce程式可以只包含map階段,
什么場景下會只需要map階段呢?
當資料只需要進行普通的過濾、決議等操作,不需要進行聚合,這個時候就不需要使用reduce階段了,在代碼層面該如何設定呢?
很簡單,在組裝Job的時候設定reduce的task數目為0就可以了,并且Reduce代碼也不需要寫了,
1、代碼如下:
package com.imooc.mr;
import org.apache.hadoop.conf.Configuration;
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.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
/**
* 只有Map階段,不包含Reduce階段
*/
public class WordCountJobNoReduce {
/**
* 組裝job=map+reduce
* 但這里沒有reduce,注意
* @param args
*/
public static void main(String[] args) {
try {
if (args.length != 2){
//如果傳遞的引數不夠,程式直接退出
System.exit(100);
}
//job需要的配置引數
Configuration conf = new Configuration();
//創建一個Job
Job job = Job.getInstance(conf);
//注意:這一行必須設定,否則在集群中執行的是找不到WordCountJob這個類
job.setJarByClass(WordCountJobNoReduce.class);
//指定輸入路徑(可以是檔案,也可以是目錄)
FileInputFormat.setInputPaths(job,new Path(args[0]));
//指定輸出路徑(只能指定一個不存在的目錄)
FileOutputFormat.setOutputPath(job,new Path(args[1]));
//指定map相關的代碼
job.setMapperClass(MyMapper.class);
//指定k2的型別
job.setMapOutputKeyClass(Text.class);
//指定v2的型別
job.setMapOutputValueClass(LongWritable.class);
//禁用reduce階段
job.setNumReduceTasks(0);
//提交job
job.waitForCompletion(true);
}catch (Exception e){
e.printStackTrace();
}
}
public static class MyMapper extends Mapper<LongWritable,Text,Text,LongWritable>{
Logger logger = LoggerFactory.getLogger(MyMapper.class);
/**
* 需要實作map函式
* 這個map函式就是可以接收k1,v1, 產生k2,v2
* @param k1
* @param v1
* @param context
* @throws IOException
* @throws InterruptedException
*/
@Override
protected void map(LongWritable k1, Text v1, Context context) throws IOException, InterruptedException {
//輸出k1,v1的值
//System.out.println("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
logger.info("<k1,v1>=<"+k1.get()+","+v1.toString()+">");
// k1代表的是每一行的行首偏移量,v1代表的是每一行內容
// 對獲取到的每一行資料進行切割,把單詞切割出來
String[] words = v1.toString().split(" ");
//迭代出來的單詞資料
for(String word:words){
//把迭代出來的單詞封裝成<k2,v2>的形式
Text k2 = new Text(word);
LongWritable v2 = new LongWritable(1L);
//把<k2,v2>寫出去
context.write(k2,v2);
}
}
}
}
2、打包上傳
打包
mvn clean package -DskipTests

備份下jar包
mv db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar_03
上傳jar包

3、執行任務
hadoop jar db_hadoop-1.0-SNAPSHOT-jar-with-dependencies.jar com.imooc.mr.WordCountJobNoReduce /test/hello.txt /out4

這里發現map執行到100%以后任務就執行成功了,reduce還是0%,因為就沒有reduce階段了,

4、查看結果
查看輸出結果,注意,這里的檔案名就是part-m-00000了
hdfs dfs -ls /out4
hdfs dfs -cat /out4/_SUCCESS
hdfs dfs -cat /out4/part-m-00000

5、web頁面查看

reduce任務為0

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436404.html
標籤:其他
上一篇:Rabbitmq的一些筆記
