Hadoop MapReduce 예제: 첫 번째 Java 프로그램 Code
⚡ 스마트 요약
Hadoop MapReduce 프로그램은 세 가지로 작성됩니다. Java 클래스, 매퍼, 리듀서 및 드라이버는 컴파일되어 JAR 파일로 패키징된 후 클러스터에 제출되어 국가별 판매량을 집계합니다.
이 튜토리얼에서는 MapReduce 예제와 함께 Hadoop을 사용하는 방법을 배웁니다. 사용된 입력 데이터는 SalesJan2009.csv이 데이터에는 제품명, 가격, 결제 방식, 고객의 도시 및 국가와 같은 판매 관련 정보가 포함되어 있습니다. 목표는 각 국가별 판매된 제품 수를 파악하는 것입니다.
최초의 Hadoop MapReduce 프로그램
이제 이것에서 맵리듀스 튜토리얼, 우리는 첫 번째를 만들 것입니다 Java 맵리듀스 프로그램:
아래 스크린샷은 2009년 1월 매출의 원본 데이터를 보여줍니다. 각 행은 하나의 거래를 나타내며, 여덟 번째 쉼표로 구분된 열에 국가 정보가 있습니다.
하둡이 설치되어 있는지 확인하십시오. 실제 작업을 시작하기 전에 사용자를 'hduser'로 변경하십시오(하둡 구성 시 사용한 ID입니다. 필요에 따라 실제 하둡 구성에 사용한 사용자 ID로 변경할 수 있습니다).
su - hduser_
아래와 같이 프롬프트가 hduser 계정으로 변경됩니다.
1단계) 프로젝트 디렉토리 및 소스 파일 생성
아래 MapReduce 예제에 나와 있는 것처럼 MapReduceTutorial이라는 이름으로 새 디렉토리를 만드세요.
sudo mkdir MapReduceTutorial
권한 부여
sudo chmod -R 777 MapReduceTutorial
세 가지를 만드세요 Java MapReduceTutorial 내부에 있는 소스 파일들을 참고하세요. 세 파일 모두 이전 버전을 사용하고 있습니다. org.apache.hadoop.mapred 이 API는 여전히 Hadoop 3.x 버전과 함께 제공됩니다.
SalesMapper.java
package SalesCountry; import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapred.*; public class SalesMapper extends MapReduceBase implements Mapper <LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); public void map(LongWritable key, Text value, OutputCollector <Text, IntWritable> output, Reporter reporter) throws IOException { String valueString = value.toString(); String[] SingleCountryData = valueString.split(","); output.collect(new Text(SingleCountryData[7]), one); } }
SalesCountryReducer.java
package SalesCountry; import java.io.IOException; import java.util.*; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapred.*; public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> { public void reduce(Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException { Text key = t_key; int frequencyForCountry = 0; while (values.hasNext()) { // replace type of value with the actual type of our value IntWritable value = (IntWritable) values.next(); frequencyForCountry += value.get(); } output.collect(key, new IntWritable(frequencyForCountry)); } }
SalesCountryDriver.java
package SalesCountry; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.*; import org.apache.hadoop.mapred.*; public class SalesCountryDriver { public static void main(String[] args) { JobClient my_client = new JobClient(); // Create a configuration object for the job JobConf job_conf = new JobConf(SalesCountryDriver.class); // Set a name of the Job job_conf.setJobName("SalePerCountry"); // Specify data type of output key and value job_conf.setOutputKeyClass(Text.class); job_conf.setOutputValueClass(IntWritable.class); // Specify names of Mapper and Reducer Class job_conf.setMapperClass(SalesCountry.SalesMapper.class); job_conf.setReducerClass(SalesCountry.SalesCountryReducer.class); // Specify formats of the data type of Input and output job_conf.setInputFormat(TextInputFormat.class); job_conf.setOutputFormat(TextOutputFormat.class); // Set input and output directories using command line arguments, //arg[0] = name of input directory on HDFS, and arg[1] = name of output directory to be created to store the output file. FileInputFormat.setInputPaths(job_conf, new Path(args[0])); FileOutputFormat.setOutputPath(job_conf, new Path(args[1])); my_client.setConf(job_conf); try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); } } }
압축을 풀면 여기에 표시된 것과 같은 세 개의 원본 파일이 생성됩니다.
이 모든 파일의 파일 권한을 확인하십시오
'읽기' 권한이 없으면 권한을 부여하세요.
2단계) 하둡 클래스패스를 내보냅니다.
아래 Hadoop 예제에 표시된 대로 클래스패스를 내보내십시오.
export CLASSPATH="$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-core-2.2.0.jar:$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-common-2.2.0.jar:$HADOOP_HOME/share/hadoop/common/hadoop-common-2.2.0.jar:~/MapReduceTutorial/SalesCountry/*:$HADOOP_HOME/lib/*"
내보낸 클래스 경로가 프롬프트에 다시 표시됩니다.
3단계) 컴파일 Java 파일
컴파일하세요 Java 해당 파일들은 Final-MapReduceHandsOn 디렉토리에 있습니다. 클래스 파일들은 패키지 디렉토리에 배치될 예정입니다.
javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java
이 경고는 무시해도 안전합니다. 단지 mapred API가 더 이상 사용되지 않는다는 것을 알려주는 것뿐입니다.
이 컴파일을 통해 현재 디렉터리에 지정된 패키지 이름으로 된 디렉터리가 생성됩니다. Java 소스 파일(예: 저희의 경우 SalesCountry)을 선택하고 컴파일된 모든 클래스 파일을 그 안에 넣으세요.
4단계) 매니페스트 파일 생성
Manifest.txt라는 새 파일을 만드세요.
sudo gedit Manifest.txt
다음 줄을 추가하세요:
Main-Class: SalesCountry.SalesCountryDriver
SalesCountry.SalesCountryDriver는 메인 클래스의 이름입니다. 이 줄을 입력한 후에는 반드시 엔터 키를 눌러야 합니다.
5단계) 내용물을 병에 담으세요
Jar 파일 만들기
jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class
jar 파일이 생성되었는지 확인
6단계) 하둡 시작하기
하둡 시작
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh
7단계) 입력 파일을 HDFS에 복사합니다.
SalesJan2009.csv 파일을 ~/inputMapReduce 폴더에 복사하세요.
이제 아래 명령어를 사용하여 ~/inputMapReduce 디렉토리를 HDFS로 복사하십시오.
$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /
우리는 이 경고를 무시해도 됩니다.
실제로 파일이 복사되었는지 확인합니다.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
8단계) MapReduce 작업을 실행합니다.
MapReduce 작업 실행
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
그러면 mapreduce_output_sales라는 출력 디렉터리가 생성됩니다. HDFS. 이 디렉토리의 내용은 국가별 제품 판매량을 포함하는 파일입니다.
9단계) 결과를 읽으세요
그 결과는 명령 인터페이스를 통해 다음과 같이 확인할 수 있습니다.
$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000
결과는 웹 인터페이스를 통해 다음과 같이 볼 수도 있습니다.
엽니다 http://localhost:50070/ 웹 브라우저에서. Hadoop 3.x 버전에서 NameNode 웹 UI는 포트로 이동했습니다. 9870, 그래서 사용하다 http://localhost:9870/ 거기 대신에.
이제 '파일 시스템 찾아보기'를 선택하고 /mapreduce_output_sales 경로로 이동하세요.
파트-r-00000 열기
SalesMapper 클래스 설명
작업이 처음부터 끝까지 실행되는 동안 다음 세 섹션에서는 각 클래스가 실제로 무엇을 하는지 자세히 설명합니다.
이 섹션에서는 SalesMapper 클래스의 구현 방식을 살펴보겠습니다.
1. 먼저 클래스에 사용할 패키지 이름을 지정합니다. 패키지 이름은 SalesCountry입니다. 컴파일 결과물인 SalesMapper.class 파일은 이 패키지 이름과 동일한 이름의 디렉터리(SalesCountry)에 저장됩니다.
그런 다음 라이브러리 패키지를 가져옵니다.
아래 스냅샷은 SalesMapper 클래스의 구현 예시를 보여줍니다.
견본 Code 설명 :
1. SalesMapper 클래스 정의-
public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {
모든 매퍼 클래스는 MapReduceBase 클래스를 상속받아야 하며, Mapper 인터페이스를 구현해야 합니다.
2. '맵' 기능 정의 -
public void map(LongWritable key, Text value, OutputCollector<Text, IntWritable> output, Reporter reporter) throws IOException
Mapper 클래스의 핵심 부분은 네 개의 인수를 받는 'map()' 메서드입니다.
'map()' 메서드를 호출할 때마다 키-값 쌍(이 코드에서는 'key'와 'value')이 전달됩니다.
'map()' 메서드는 인수로 받은 입력 텍스트를 분할하는 것으로 시작합니다. 각 줄을 필드로 분할합니다.
String valueString = value.toString(); String[] SingleCountryData = valueString.split(",");
여기서 ','는 구분 기호로 사용됩니다.
이후, 'SingleCountryData' 배열의 7번째 인덱스에 있는 레코드와 값 '1'을 사용하여 쌍이 형성됩니다.
output.collect(new Text(SingleCountryData[7]), one);
'SingleCountryData' 배열에서 국가 데이터가 7번째 인덱스에 위치해 있기 때문에 해당 레코드를 선택합니다.
입력 데이터는 아래와 같은 형식입니다(국가는 7번째 인덱스이며, 0부터 시작합니다).
Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude
매퍼의 출력은 'OutputCollector'의 'collect()' 메서드를 사용하여 출력되는 키-값 쌍입니다.
SalesCountryReducer 클래스 설명
이 섹션에서는 SalesCountryReducer 클래스의 구현 방식을 살펴보겠습니다.
1. 먼저 클래스에 사용할 패키지 이름을 지정합니다. 패키지 이름은 SalesCountry입니다. 컴파일 결과물인 SalesCountryReducer.class 파일은 이 패키지 이름과 동일한 이름의 디렉터리(SalesCountry)에 저장됩니다.
그런 다음 라이브러리 패키지를 가져옵니다.
아래 스냅샷은 SalesCountryReducer 클래스의 구현 예시를 보여줍니다.
Code 설명 :
1. SalesCountryReducer 클래스 정의-
public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {
여기서 처음 두 데이터 유형인 'Text'와 'IntWritable'은 리듀서에 입력되는 키-값의 데이터 유형입니다.
매퍼의 출력은 다음과 같은 형식입니다. , 매퍼의 출력은 리듀서의 입력이 됩니다. 따라서 데이터 유형과 일치하도록 여기서는 Text와 IntWritable을 데이터 유형으로 사용합니다.
마지막 두 데이터 유형인 'Text'와 'IntWritable'은 리듀서가 키-값 쌍 형태로 생성하는 출력 데이터 유형입니다.
모든 리듀서 클래스는 MapReduceBase 클래스를 상속받아야 하며, Reducer 인터페이스를 구현해야 합니다.
2. '감소' 기능 정의 -
public void reduce( Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
reduce() 메서드의 입력은 여러 값으로 이루어진 리스트를 가진 키입니다.
예를 들어, 우리의 경우에는 다음과 같습니다.
, , , , , .
이것은 감속기에 다음과 같이 주어집니다.
따라서 이러한 형식의 인수를 받아들이기 위해서는 먼저 Text와 Iterator라는 두 가지 데이터 유형이 사용됩니다. 텍스트는 키와 이터레이터의 데이터 유형입니다. 해당 키에 대한 값 목록을 나타내는 데이터 유형입니다.
다음 인수는 OutputCollector 유형입니다. 이는 감속기 단계의 출력을 수집합니다.
reduce() 메서드는 키 값을 복사하고 빈도수를 0으로 초기화하는 것으로 시작합니다.
Text key = t_key;int frequencyForCountry = 0;
그런 다음 'while' 루프를 사용하여 키와 연결된 값 목록을 순회하고 모든 값을 합산하여 최종 빈도를 계산합니다.
while (values.hasNext()) { // replace type of value with the actual type of our value IntWritable value = (IntWritable) values.next(); frequencyForCountry += value.get(); }
이제 키와 얻어진 빈도수 형태로 결과를 출력 수집기로 보냅니다.
아래 코드는 이를 수행합니다.
output.collect(key, new IntWritable(frequencyForCountry));
SalesCountryDriver 클래스 설명
이 섹션에서는 SalesCountryDriver 클래스의 구현 방식을 살펴보겠습니다.
1. 먼저 클래스에 사용할 패키지 이름을 지정합니다. 패키지 이름은 SalesCountry입니다. 컴파일 결과물인 SalesCountryDriver.class 파일은 이 패키지 이름과 동일한 이름의 디렉터리(SalesCountry)에 저장됩니다.
다음은 패키지 이름을 지정하는 줄과 라이브러리 패키지를 가져오는 코드입니다.
2. 새 클라이언트 작업, 구성 개체를 생성하고 매퍼 및 감속기 클래스를 광고할 드라이버 클래스를 정의합니다.
드라이버 클래스는 MapReduce 작업이 실행되도록 설정하는 역할을 담당합니다. 하둡이 클래스에서는 작업 이름, 입력/출력 데이터 유형, 매퍼 및 리듀서 클래스의 이름을 지정합니다.
3. 아래 코드 조각에서는 각각 입력 데이터 세트를 사용하고 출력을 생성하는 데 사용되는 입력 및 출력 디렉터리를 설정합니다.
arg[0] 및 arg[1]은 MapReduce 실습에서 제공된 명령과 함께 전달되는 명령줄 인수입니다. 즉,
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
4. 작업을 시작하세요
아래 코드는 MapReduce 작업의 실행을 시작합니다.
try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); }





















