Hadoop MapReduce 예제: 첫 번째 Java 프로그램 Code

⚡ 스마트 요약

Hadoop MapReduce 프로그램은 세 가지로 작성됩니다. Java 클래스, 매퍼, 리듀서 및 드라이버는 컴파일되어 JAR 파일로 패키징된 후 클러스터에 제출되어 국가별 판매량을 집계합니다.

  • 🔘 데이터 세트 : SalesJan2009.csv 파일에는 각 줄마다 하나의 거래 내역이 있으며, 여덟 번째 열(쉼표로 구분됨)에 국가 정보가 있습니다.
  • ☑️ 매퍼: SalesMapper는 각 행을 분할하고 국가와 상수 값 1을 쌍으로 출력합니다.
  • 감속기 : SalesCountryReducer는 각 국가별로 도착하는 물량을 합산하여 단일 빈도수로 계산합니다.
  • 🧪 드라이버 : SalesCountryDriver 클래스는 작업 이름을 지정하고, 키 및 값 유형을 선언하며, 두 클래스를 연결합니다.
  • 🛠️ 빌드 : javac로 컴파일하고, Main-Class 매니페스트 항목을 추가한 다음, jar cfm으로 모든 것을 패키징하세요.
  • ⚠️ 운영: CSV 파일을 HDFS에 복사하고, jar 파일을 제출한 다음, 출력 디렉토리에서 part-00000 부분을 읽으십시오.

Hadoop MapReduce 예제를 통해 첫 번째 생성을 진행합니다. Java 프로그램

이 튜토리얼에서는 MapReduce 예제와 함께 Hadoop을 사용하는 방법을 배웁니다. 사용된 입력 데이터는 SalesJan2009.csv이 데이터에는 제품명, 가격, 결제 방식, 고객의 도시 및 국가와 같은 판매 관련 정보가 포함되어 있습니다. 목표는 각 국가별 판매된 제품 수를 파악하는 것입니다.

최초의 Hadoop MapReduce 프로그램

이제 이것에서 맵리듀스 튜토리얼, 우리는 첫 번째를 만들 것입니다 Java 맵리듀스 프로그램:

아래 스크린샷은 2009년 1월 매출의 원본 데이터를 보여줍니다. 각 행은 하나의 거래를 나타내며, 여덟 번째 쉼표로 구분된 열에 국가 정보가 있습니다.

2009년 1월 매출 데이터가 거래 내역 열이 표시된 스프레드시트로 열렸습니다.

하둡이 설치되어 있는지 확인하십시오. 실제 작업을 시작하기 전에 사용자를 'hduser'로 변경하십시오(하둡 구성 시 사용한 ID입니다. 필요에 따라 실제 하둡 구성에 사용한 사용자 ID로 변경할 수 있습니다).

su - hduser_

아래와 같이 프롬프트가 hduser 계정으로 변경됩니다.

hduser 계정으로 전환 후 터미널 프롬프트

1단계) 프로젝트 디렉토리 및 소스 파일 생성

아래 MapReduce 예제에 나와 있는 것처럼 MapReduceTutorial이라는 이름으로 새 디렉토리를 만드세요.

sudo mkdir MapReduceTutorial

mkdir 명령어를 사용하여 MapReduceTutorial 디렉토리를 생성합니다.

권한 부여

sudo chmod -R 777 MapReduceTutorial

chmod 명령어를 사용하여 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();
        }
    }
}

여기서 파일을 다운로드하세요

압축을 풀면 여기에 표시된 것과 같은 세 개의 원본 파일이 생성됩니다.

Extracted archive에는 SalesMapper, SalesCountryReducer 및 SalesCountryDriver 목록이 있습니다.

이 모든 파일의 파일 권한을 확인하십시오

세 가지 항목의 파일 권한을 보여주는 긴 목록 Java 소스 파일

'읽기' 권한이 없으면 권한을 부여하세요.

chmod 명령어를 사용하여 읽기 권한을 추가합니다. Java 소스 파일

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/*"

내보낸 클래스 경로가 프롬프트에 다시 표시됩니다.

Hadoop CLASSPATH 변수를 내보낸 후 셸 프롬프트가 표시됩니다.

3단계) ​​컴파일 Java 파일

컴파일하세요 Java 해당 파일들은 Final-MapReduceHandsOn 디렉토리에 있습니다. 클래스 파일들은 패키지 디렉토리에 배치될 예정입니다.

javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java

javac 컴파일 출력에서 ​​사용 중단 경고가 표시됩니다.

이 경고는 무시해도 안전합니다. 단지 mapred API가 더 이상 사용되지 않는다는 것을 알려주는 것뿐입니다.

이 컴파일을 통해 현재 디렉터리에 지정된 패키지 이름으로 된 디렉터리가 생성됩니다. Java 소스 파일(예: 저희의 경우 SalesCountry)을 선택하고 컴파일된 모든 클래스 파일을 그 안에 넣으세요.

SalesCountry 패키지 디렉터리에는 컴파일된 클래스 파일 세 개가 포함되어 있습니다.

4단계) 매니페스트 파일 생성

Manifest.txt라는 새 파일을 만드세요.

sudo gedit Manifest.txt

다음 줄을 추가하세요:

Main-Class: SalesCountry.SalesCountryDriver

gedit Manifest.txt의 Main-Class 항목을 보여주는 창

SalesCountry.SalesCountryDriver는 메인 클래스의 이름입니다. 이 줄을 입력한 후에는 반드시 엔터 키를 눌러야 합니다.

5단계) 내용물을 병에 담으세요

Jar 파일 만들기

jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class

jar 명령어를 사용하여 매니페스트에서 ProductSalePerCountry.jar 파일을 빌드합니다.

jar 파일이 생성되었는지 확인

ProductSalePerCountry.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 /

copyFromLocal 명령의 출력에 네이티브 라이브러리 경고가 표시됩니다.

우리는 이 경고를 무시해도 됩니다.

실제로 파일이 복사되었는지 확인합니다.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls 명령을 실행하면 inputMapReduce 디렉터리 안에 SalesJan2009.csv 파일이 나열됩니다.

8단계) MapReduce 작업을 실행합니다.

MapReduce 작업 실행

$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales

클러스터에서 SalePerCountry 작업이 실행되는 동안 콘솔에 출력되는 내용입니다.

그러면 mapreduce_output_sales라는 출력 디렉터리가 생성됩니다. HDFS. 이 디렉토리의 내용은 국가별 제품 판매량을 포함하는 파일입니다.

9단계) 결과를 읽으세요

그 결과는 명령 인터페이스를 통해 다음과 같이 확인할 수 있습니다.

$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000

hdfs dfs -cat 명령으로 출력되는 국가 및 판매량 쌍

결과는 웹 인터페이스를 통해 다음과 같이 볼 수도 있습니다.

엽니다 http://localhost:50070/ 웹 브라우저에서. Hadoop 3.x 버전에서 NameNode 웹 UI는 포트로 이동했습니다. 9870, 그래서 사용하다 http://localhost:9870/ 거기 대신에.

브라우저에서 보는 NameNode 웹 인터페이스 홈페이지

이제 '파일 시스템 찾아보기'를 선택하고 /mapreduce_output_sales 경로로 이동하세요.

HDFS 파일 브라우저에서 mapreduce_output_sales 디렉토리를 보여줍니다.

파트-r-00000 열기

part-r-00000 결과 파일이 HDFS 브라우저에서 열렸습니다.

SalesMapper 클래스 설명

작업이 처음부터 끝까지 실행되는 동안 다음 세 섹션에서는 각 클래스가 실제로 무엇을 하는지 자세히 설명합니다.

이 섹션에서는 SalesMapper 클래스의 구현 방식을 살펴보겠습니다.

1. 먼저 클래스에 사용할 패키지 이름을 지정합니다. 패키지 이름은 SalesCountry입니다. 컴파일 결과물인 SalesMapper.class 파일은 이 패키지 이름과 동일한 이름의 디렉터리(SalesCountry)에 저장됩니다.

그런 다음 라이브러리 패키지를 가져옵니다.

아래 스냅샷은 SalesMapper 클래스의 구현 예시를 보여줍니다.

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 클래스의 구현 예시를 보여줍니다.

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)에 저장됩니다.

다음은 패키지 이름을 지정하는 줄과 라이브러리 패키지를 가져오는 코드입니다.

SalesCountryDriver의 패키지 선언 및 가져오기 명세서

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();
}

자주 묻는 질문

세 클래스 모두 원래 MapReduce API인 org.apache.hadoop.mapred를 임포트합니다. Hadoop 3.x 버전에서는 여전히 org.apache.hadoop.mapred를 포함하고 실행하지만, 새로운 작업은 일반적으로 JobConf와 JobClient를 Configuration과 Job으로 대체하는 org.apache.hadoop.mapreduce를 사용하여 작성됩니다.

과거 작업 이력을 기반으로 학습된 모델은 실행 시간을 예측하고, 분할 크기와 리듀서 개수를 제안하며, 실행이 완료되기 전에 편차를 감지합니다. 또한, 스필된 레코드 수나 실패한 작업 수가 정상 범위를 벗어나는 작업을 표시합니다.

Copilot은 클래스 시그니처, 제네릭, 임포트 및 드라이버 구성 호출과 같은 반복적인 부분을 잘 처리합니다. 하지만 국가 정보를 저장하는 열과 같은 비즈니스 로직은 개발자가 실제 스키마와 비교하여 확인해야 합니다.

Hadoop은 기존 출력 디렉터리에 쓰기를 거부하여 완료된 결과가 덮어쓰여지지 않습니다. `hdfs dfs -rm -r` 명령어를 사용하여 `mapreduce_output_sales` 디렉터리를 재귀적으로 삭제하거나, 다음 실행 시 다른 출력 경로를 지정하십시오.

jar 도구는 줄 바꿈 문자가 없는 마지막 매니페스트 줄을 무시하므로 Main-Class가 아무런 경고 없이 삭제됩니다. 그 결과 jar 파일은 오류 없이 빌드되지만, 메인 클래스가 기록되지 않아 런타임 오류가 발생합니다.

네. 예제에서는 모든 jar 파일 이름에 2.2.0 버전이 포함되어 있습니다. 다른 릴리스에서는 해당 버전으로 바꾸거나, 설치된 배포판에 필요한 모든 jar 파일을 출력하는 `hadoop classpath` 명령어를 사용하면 됩니다.

아니요. 단일 노드 유사 분산 설치에서는 해당 명령어가 HDFS와 YARN이 시작되었다고 가정하기 때문에 변경 없이 실행됩니다. 동일한 jar 파일이 코드 변경 없이 실제 다중 노드 클러스터에도 배포됩니다.

매퍼에서 배열 인덱스를 7에서 5로 변경하십시오. City는 SalesJan2009.csv 파일의 여섯 번째 열입니다. 다시 컴파일하고, jar 파일을 다시 빌드한 다음, 새 출력 디렉터리를 대상으로 작업을 실행하십시오.

이 게시물을 요약하면 다음과 같습니다.