Amazon Ads

顯示具有 Hadoop 標籤的文章。 顯示所有文章
顯示具有 Hadoop 標籤的文章。 顯示所有文章

2015年6月3日 星期三

【筆記】使用MapReduce來替換字串(replace string)的簡單程式

拜歐先把要替換的檔案上傳到HDFS中。

上傳完成後,使用下列指令檢視其內容:
hdfs dfs -cat /user/javakid/replace/replace.txt
輸入的檔案為replace.txt,其內容如下:
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
foo1234foo567foo890foo123foo
再來開始寫程式,拜歐是用Intellij IDEA來撰寫程式,在專案設定中,將編譯的輸出目錄(compile output path)設定為/home/javakid/study/hadoop/classpath。

第一個要寫的是 mapper:
package idv.jk.study.hadoop.myself;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

import java.io.IOException;

/**
 * Created by javakid on 2015/6/3.
 */
public class StringReplaceMapper
        extends Mapper<LongWritable, Text, NullWritable, Text>
{
    private static final String TARGET = "foo";
    private static final String REPLACEMENT = "bar";

    @Override
    protected void map(LongWritable key, Text value, Context context)
            throws IOException, InterruptedException
    {
        String line = value.toString();

        if(line.indexOf(TARGET) >= 0)
        {
            context.write(NullWritable.get(),
                            new Text(line.replaceAll(TARGET, REPLACEMENT)));
        }
    }
}
在mapper的輸出,拜歐要的只是替換後的每一行字串,這些值的key是什麼,並不重要,所以在上列程式中的第27行,用NullWritable來做為輸出的key。

再來要寫的是 reducer:
package idv.jk.study.hadoop.myself;

import org.apache.hadoop.io.NullWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;

import java.io.IOException;

/**
 * Created by javakid on 2015/6/3.
 */
public class StringReplaceReducer
        extends Reducer<NullWritable, Text, NullWritable, Text>
{
    @Override
    protected void reduce(NullWritable key, Iterable<Text> values, Context context)
            throws IOException, InterruptedException
    {
        for(Text text : values)
        {
            context.write(NullWritable.get(), text);
        }
    }
}
在reducer中,拜歐要的也只是將已經在mapper替換好的字串輸出,在第21行一樣使用NullWritable來做為輸出的key。

最後是主程式:
package idv.jk.study.hadoop.myself;

import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.NullWritable;
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;

/**
 * Created by javakid on 2015/6/3.
 */
public class StringReplace
{
    public static void main(String[] argv)
            throws IOException, ClassNotFoundException, InterruptedException
    {
        if(argv.length != 2)
        {
            System.err.println("Usage: StringReplace <input path> <output path>");
            System.exit(-1);
        }

        Job job = new Job();
        job.setJarByClass(StringReplace.class);
        job.setJobName("String replacement");

        FileInputFormat.addInputPath(job, new Path(argv[0]));
        FileOutputFormat.setOutputPath(job, new Path(argv[1]));

        job.setMapperClass(StringReplaceMapper.class);
        job.setReducerClass(StringReplaceReducer.class);

        job.setOutputKeyClass(NullWritable.class);
        job.setOutputValueClass(Text.class);

        System.out.println(job.waitForCompletion(true) ? 0 : 1);
    }
}
上列setOutputKeyClass和setOutputValueClass這兩個方法控制reduce方法中輸出的型別,而且與reducer中產出的輸出中設定,兩者必須是相同的。

先將HADOOP_CLASSPATH設定好後,用hadoop指令來執行主程式:
export HADOOP_CLASSPATH=/home/javakid/study/hadoop/classpath
hadoop idv.jk.study.hadoop.myself.StringReplace \
/user/javakid/replace/replace.txt /user/javakid/replace/output
執行完成後,使用下列指令來確認結果:
hdfs dfs -cat /user/javakid/replace/output/part-r-00000
預期結果如下:
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
bar1234bar567bar890bar123bar
你可以在這裡找到原始碼。

2015年5月30日 星期六

【筆記】在Ubuntu 14.04建立Hadoop多節點叢集-分散式架構(Multi-node cluster)

這裡使用的版本是2.6.0,之前有些設定如設定JAVA_HOME
或是ssh等,在Setting up a Single Node Cluster己經完成,若有撞牆的地方,可以回去參考一下。

在拜歐的cluster中,會有四台機器,分別為:
  • javakid01:做為master,會用它來跑NameNode和SecondaryNameNode
  • jkserver01:做為slaves,會用它來跑DataNode
  • jkserver02:做為slaves,會用它來跑DataNode
  • jkserver03:做為slaves,會用它來跑DataNode

若要改hostname,你可以這樣做:
sudo vi /etc/hostname
上列四台機器中,先確認已經安裝好Java和Hadoop。

先對javakid01(master),進行設定,在/etc/hosts中加入:
10.211.55.10    javakid01
10.211.55.7     jkserver01
10.211.55.13    jkserver02
10.211.55.14    jkserver03
再來編輯$HADOOP_INSTALL/etc/hadoop/core-site.xml加入下列設定:

        
                fs.defaultFS
                hdfs://javakid01:9000
        

接著編輯$HADOOP_INSTALL/etc/hadoop/hdfs-site.xml加入下列設定:

        
                dfs.replication
                3
        
        
                dfs.namenode.name.dir
                file:/opt/data/hadoop/hadoop_data/hdfs/namenode
        

上面設定好後,使用下列指令建立上列設定的目錄,建立前請確定有足夠的權限:
mkdir -p opt/data/hadoop/hadoop_data/hdfs/namenode
接著編輯$HADOOP_INSTALL/etc/hadoop/mapred-site.xml加入下列設定:

        
                mapred.job.tracker
                javakid01:54311
        

接著編輯$HADOOP_INSTALL/etc/hadoop/yarn-site.xml加入下列設定:


        
                yarn.nodemanager.aux-services
                mapreduce_shuffle
        
        
                yarn.nodemanager.aux-services.mapreduce.shuffle.class
                org.apache.hadoop.mapred.ShuffleHandler
        
        
                yarn.resourcemanager.resource-tracker.address
                javakid01:8025
        
        
                yarn.resourcemanager.scheduler.address
                javakid01:8030
        
        
                yarn.resourcemanager.address
                javakid01:8050
        

最後,在$HADOOP_INSTALL/etc/hadoop/slaves中加入做為DataNode的hostname
jkserver01
jkserver02
jkserver03
下面要開始對做為slaves的機器進行設定,其中jkserver01、jkserver02與jkserver03的設定幾乎是相同的。

一樣先在/etc/hosts中加入:
10.211.55.10    javakid01
10.211.55.7     jkserver01
10.211.55.13    jkserver02
10.211.55.14    jkserver03
請注意!上列的IP會因機器不同而有差異,這裡設定是以我的機器為例。

再來編輯$HADOOP_INSTALL/etc/hadoop/core-site.xml加入下列設定:

        
                fs.defaultFS
                hdfs://javakid01:9000
        

接著編輯$HADOOP_INSTALL/etc/hadoop/hdfs-site.xml加入下列設定:

        
                dfs.replication
                3
        
        
                dfs.datanode.data.dir
                file:/opt/data/hadoop/hadoop_data/hdfs/datanode
        

上面設定好後,使用下列指令建立上列設定的目錄,建立前請確定有足夠的權限:
mkdir -p opt/data/hadoop/hadoop_data/hdfs/datanode
接著編輯$HADOOP_INSTALL/etc/hadoop/mapred-site.xml加入下列設定:

        
                mapred.job.tracker
                javakid01:54311
        

接著編輯$HADOOP_INSTALL/etc/hadoop/yarn-site.xml加入下列設定:


        
                yarn.nodemanager.aux-services
                mapreduce_shuffle
        
        
                yarn.nodemanager.aux-services.mapreduce.shuffle.class
                org.apache.hadoop.mapred.ShuffleHandler
        
        
                yarn.resourcemanager.resource-tracker.address
                javakid01:8025
        
        
                yarn.resourcemanager.scheduler.address
                javakid01:8030
        
        
                yarn.resourcemanager.address
                javakid01:8050
        

最後,也要將讓DataNode知道slaves的設定。在$HADOOP_INSTALL/etc/hadoop/slaves中加入做為DataNode的hostname
jkserver01
jkserver02
jkserver03
接下來,為了讓master不用密碼就可登入slaves的話,需要進行SSH連線的設定。

回到javakid01這台機器,在終端機下執行:
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa
cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys
然後執行:
ssh-copy-id -i ~/.ssh/id_rsa.pub javakid@jkserver01
完成後,由javakid01登入jkserver01,第一次登入會需要在jkserver01上的使用者密碼:
ssh jkserver01
對jkserver02與jkserver03進行相同設定,在javakid01執行:
ssh-copy-id -i ~/.ssh/id_rsa.pub javakid@jkserver02
登入jkserver02:
ssh jkserver02
接著執行:
ssh-copy-id -i ~/.ssh/id_rsa.pub javakid@jkserver03
登入jkserver03:
ssh jkserver03
全部設定完成後,在javakid01執行下列指令對HDFS進行格式化:
hadoop namenode -format
接著啟動 HDFS Daemons:
$HADOOP_INSTALL/sbin/start-dfs.sh
成功啟動後,在javakid01執行jps指令:
jps
應該可以看到下列的結果:
27964 NameNode
28220 SecondaryNameNode
29695 Jps
再來若是到jkserver01、jkserver02或jkserver03任一台機器執行jps,應該可以看到下列結果:
2204 DataNode
2378 Jps
若在瀏覽器輸入http://javakid01:50070這個網址,應該在Datanodes這個頁籤看到三台機器的資訊:


2015年5月26日 星期二

【筆記】Hadoop名詞解釋

名詞意思備註
metricsHDFS與MapReduce daemon會收集一些與事件以及量測相關的資訊,這些總稱為metrics例如,datanode會收集已寫入的byte數、已複製的block總數等metrics
data locality optimization Hadoop在執行map task時,會盡全力將該task放到存放輸入該task資料的那台node去執行,以免使用到寶貴的頻寬進行資料的傳送
shuffle A proccess by which the system performs the sort and transfers the map outputs to the reducers as inputs.
speculative execution Hadoop doesn't try to diagnose and fix slow-running tasks; instead, it tries to detect when a task is running slower than expected and launches another equivalent task as a backup.

2015年2月25日 星期三

【筆記】解決Hadoop在主機重開機後,namenode無法啟動的問題

在使用pseudo-distributed mode設定好Hadoop以後,執行hadoop namenode -format,再執行start-dfs.sh後,即可正常啟動Namenode等,順利把檔案放到Hadoop上去,或是瀏覽Hadoop上的檔案。

但有時只要每次重開機後,再執行start-dfs.sh啟動Hadoop服務,下jps去檢查時,會看不到Namenode,這時可能的原因是沒在core-site.xml中設定hadoop.tmp.dir。

在沒有設定hadoop.tmp.dir這個參數的情況下,Hadoop的暫存資料夾預設為/tmp/hadoop-你的使用名稱,如:


而這個資料夾在每次重開機後,即會被清除,因此會造成Namenode在重開機無法正常啟動的情況,解決方式就是在core-site.xml中設定hadoop.tmp.dir來指定Hadoop的暫存資料夾:

        hadoop.tmp.dir
 /home/javakid/study/hadoop/temp_data

設定完成後,再執行hadoop namenode -format,下次重新啟動後,再啟動DFS後,Namenode就會正常啟動。