Hadoop MapReduce 조인 및 카운터 예시

⚡ 스마트 요약

MapReduce 조인은 매퍼 내부 또는 리듀서 내부에서 공유 키를 기준으로 두 개의 대규모 데이터 세트를 결합하는 반면, MapReduce 카운터는 작업에 대한 통계를 수집하여 잘못된 레코드를 추측이 아닌 측정을 통해 찾아낼 수 있도록 합니다.

  • 🔘 가입 기본 사항: 두 데이터 세트 중 더 작은 데이터 세트가 모든 데이터 노드에 배포되어 조회용으로 사용됩니다.
  • ☑️ 지도 측 조인: 맵 함수를 실행하기 전에 각 입력값을 파티션하고, 동일하게 분할하고, 조인 키를 기준으로 정렬해야 합니다.
  • 축소측 조인: 조인 키를 공유하는 모든 튜플이 동일한 리듀서에 도달하므로 파티셔닝이 필요하지 않습니다.
  • 🧪 실제 예제: DeptName.txt와 DeptStrength.txt 파일은 HDFS에 복사되고, 패키지된 JAR 파일을 통해 Dept_ID를 기준으로 조인됩니다.
  • 🛠️ 카운터 유형: 모든 작업에는 5개의 내장 카운터 그룹이 포함되어 있으며, 사용자 정의 카운터는 다음과 같이 선언됩니다. Java 열거형.
  • ⚠️ 반대 용도: 누락되거나 유효하지 않은 레코드가 발생할 때마다 카운터를 증가시키면 데이터 품질 문제가 작업 보고서의 수치로 표시됩니다.

Hadoop MapReduce 조인 및 카운터 튜토리얼 (실행 예제 포함)

MapReduce에서 Join이란 무엇인가요?

MapReduce의 조인 연산은 두 개의 대규모 데이터셋을 결합하는 데 사용됩니다. 하지만 이 과정은 실제 조인 연산을 수행하기 위해 많은 코드를 작성해야 합니다. 두 데이터셋을 조인하는 과정은 각 데이터셋의 크기를 비교하는 것으로 시작됩니다. 한 데이터셋이 다른 데이터셋보다 작으면, 작은 데이터셋이 클러스터의 모든 데이터 노드에 분산됩니다.

일단 참여하세요 MapReduce 데이터가 분산되어 있는 경우, 매퍼 또는 리듀서는 더 작은 데이터 세트를 사용하여 더 큰 데이터 세트에서 일치하는 레코드를 찾은 다음 이러한 레코드를 결합하여 출력 레코드를 생성합니다.

조인 유형

실제 조인이 수행되는 위치에 따라 하둡의 조인은 두 가지 종류로 분류됩니다.

  1. 지도 측면 조인 — 매퍼에서 조인을 수행하는 경우, 이를 맵 측 조인이라고 합니다. 이 유형에서는 데이터가 맵 함수에 의해 실제로 사용되기 전에 조인이 수행됩니다. 각 맵의 입력은 파티션 형태여야 하며 정렬된 상태여야 합니다. 또한, 파티션의 개수는 모두 같아야 하고 조인 키를 기준으로 정렬되어 있어야 합니다.
  2. 축소 측면 연결 — 리듀서에서 조인이 수행될 때, 이를 리듀스 측 조인이라고 합니다. 이 조인에서는 데이터셋이 구조화된 형태(또는 파티셔닝된 형태)일 필요가 없습니다. 맵 측 처리는 조인 키와 두 테이블의 해당 튜플을 출력합니다. 이 처리의 결과로, 동일한 조인 키를 가진 모든 튜플이 동일한 리듀서로 전달되고, 리듀서는 동일한 조인 키를 가진 레코드들을 조인합니다.

Hadoop의 전체 조인 프로세스 흐름은 아래 다이어그램에 설명되어 있습니다.

하둡에서 맵 측 조인과 리듀스 측 조인을 비교하는 프로세스 흐름도
Hadoop MapReduce의 조인 유형

두 가지 변형이 명확해졌으므로 다음 섹션에서는 두 개의 작은 부서 파일에 대한 축소 측 조인을 살펴보겠습니다.

두 개의 데이터 세트를 결합하는 방법: MapReduce 예

두 개의 서로 다른 파일에 두 세트의 데이터가 있습니다(아래 참조). 두 파일 모두에 공통 키인 Dept_ID가 있습니다. 목표는 MapReduce Join을 사용하여 이 두 파일을 결합하는 것입니다.

첫 번째 입력 파일에는 부서 이름과 함께 부서 ID가 나열되어 있습니다.

1 파일
두 번째 입력 파일에는 부서 ID와 부서별 인원수가 함께 나열되어 있습니다.

2 파일

입력: 입력 데이터 세트는 DeptName.txt와 DeptStrength.txt라는 텍스트 파일입니다.

여기에서 입력 파일 다운로드

당신이 가지고 있는지 확인 하둡 설치가 완료되었습니다. MapReduce Join 예제의 실제 프로세스를 시작하기 전에 사용자를 'hduser'로 변경하십시오(Hadoop 구성 시 사용한 ID이며, Hadoop 구성 시 사용한 사용자 ID로 변경할 수 있습니다).

su - hduser_

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

su 명령어를 사용하여 hduser 계정으로 전환한 후의 터미널 화면

단계 1) zip 파일을 원하는 위치에 복사하세요.

다운로드한 MapReduceJoin 아카이브가 선택한 작업 디렉터리에 배치되었습니다.

단계 2) Zip 파일 압축 풀기

sudo tar -xvf MapReduceJoin.tar.gz

전tractar가 압축 파일을 풀면서 파일 이름들이 스크롤되어 지나갑니다.

콘솔에 파일 목록 표시 (예: ...)tracMapReduceJoin.tar.gz에서 가져온 ted

단계 3) MapReduceJoin/ 디렉터리로 이동합니다.

cd MapReduceJoin/

MapReduceJoin 디렉토리로 이동한 후의 셸 프롬프트

단계 4) 하둡 시작

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

두 스크립트 모두 실행하는 데몬 목록을 출력합니다.

HDFS 및 YARN 데몬 스크립트의 시작 메시지

단계 5) DeptStrength.txt 및 DeptName.txt는 이 MapReduce Join 예제 프로그램에 사용되는 입력 파일입니다.

이 파일들을 다음 위치로 복사해야 합니다. HDFS 아래 명령어를 사용하세요.

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal DeptStrength.txt DeptName.txt /

두 입력 텍스트 파일 모두 HDFS 루트 디렉터리에 복사되었습니다.

단계 6) 아래 명령을 사용하여 프로그램을 실행하십시오.

$HADOOP_HOME/bin/hadoop jar MapReduceJoin.jar MapReduceJoin/JoinDriver/DeptStrength.txt /DeptName.txt /output_mapreducejoin

먼저 명령어가 출력되고, 그 후 작업이 콘솔에 진행 상황을 보고합니다.

명령줄에서 패키지된 MapReduceJoin jar 파일을 실행합니다.

콘솔 출력 tracMapReduce 조인 작업의 진행 상황을 확인하세요.

단계 7) 실행 후 출력 파일('part-00000'이라는 이름)은 HDFS의 /output_mapreducejoin 디렉터리에 저장됩니다.

명령줄 인터페이스를 사용하여 결과를 볼 수 있습니다.

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

cat 명령어를 사용하여 HDFS에서 출력된 입사 부서 기록

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

파일 시스템 브라우저에 접근하기 위해 사용되는 하둡 웹 인터페이스 랜딩 페이지

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

HDFS 파일 시스템 뷰에서 output_mapreducejoin 디렉터리를 탐색합니다.

파트-r-00000 열기

브라우저 보기에서 part-r-00000 출력 파일을 선택합니다.

결과가 표시됩니다

브라우저에 표시되는 부서명 및 부서 규모 행

알림: 다음에 이 프로그램을 실행하기 전에 출력 디렉터리 /output_mapreducejoin을 삭제해야 합니다.

$HADOOP_HOME/bin/hdfs dfs -rm -r /output_mapreducejoin

대안은 출력 디렉터리에 다른 이름을 사용하는 것입니다.

조인은 데이터의 형태를 알려줍니다. 다음에 다룰 카운터는 데이터를 생성한 작업의 동작 방식을 알려줍니다.

MapReduce의 카운터란 무엇입니까?

MapReduce에서 카운터는 MapReduce 작업 및 이벤트에 대한 통계 정보를 수집하고 측정하는 데 사용되는 메커니즘입니다. 카운터는 다음과 같은 정보를 유지합니다. tracMapReduce에서 발생하는 연산 횟수 및 연산 진행률과 같은 다양한 작업 통계에 대한 카운터입니다. 카운터는 MapReduce에서 문제 진단에 사용됩니다.

Hadoop 카운터는 맵이나 축소를 위한 코드에 로그 메시지를 넣는 것과 유사합니다. 이 정보는 MapReduce 작업 처리 시 문제를 진단하는 데 유용할 수 있습니다.

일반적으로 하둡의 이러한 카운터는 프로그램(맵 또는 리듀스)에 정의되며, 특정 이벤트나 조건(해당 카운터에 특정한 조건)이 발생할 때 실행 중에 증가합니다. 하둡 카운터의 매우 유용한 활용 사례는 다음과 같습니다. trac입력 데이터 세트에서 유효한 레코드와 유효하지 않은 레코드를 각각 k개 추출합니다.

MapReduce 카운터 유형

MapReduce 카운터에는 기본적으로 두 가지 유형이 있습니다.

  1. Hadoop 내장 카운터: 작업별로 존재하는 내장형 Hadoop 카운터가 있습니다. 다음은 내장된 카운터 그룹입니다.
    • MapReduce 작업 카운터 — 실행 시간 동안 작업별 정보(예: 입력 레코드 수)를 수집합니다.
    • 파일 시스템 카운터 — 작업에서 읽거나 쓴 바이트 수와 같은 정보를 수집합니다.
    • FileInputFormat 카운터 — FileInputFormat을 통해 읽은 바이트 수에 대한 정보를 수집합니다.
    • FileOutputFormat 카운터 — FileOutputFormat을 통해 기록된 바이트 수에 대한 정보를 수집합니다.
    • 작업 카운터 — 이 카운터는 작업에 대해 시작된 작업 수와 같은 작업 전체 통계를 기록합니다.
  2. 사용자 정의 카운터: 내장 카운터 외에도 사용자는 프로그래밍 언어에서 제공하는 유사한 기능을 사용하여 자체 카운터를 정의할 수 있습니다. 예를 들어, Java'enum'은 사용자 정의 카운터를 정의하는 데 사용됩니다.

💡 버전 참고: 구인 현황은 Job에서 관리했습니다.TracMRv1에서는 ker 역할을 합니다. YARN에서는 해당 역할이 MapReduce ApplicationMaster에 속하므로 카운터 이름은 유지되지만 이를 보고하는 구성 요소가 변경됩니다.

작업은 무제한의 카운터를 선언할 수 없습니다. mapreduce.job.counters.max 이 설정은 기본적으로 작업당 총 개수를 120개로 제한하며, 그 이상을 선언하는 작업은 오류를 발생시킵니다. LimitExceededException따라서 카운터는 키별 집계보다는 몇 개의 신호를 종합적으로 집계하기 위한 것입니다.

카운터 예

누락되거나 유효하지 않은 값의 개수를 세는 카운터가 포함된 MapClass의 예입니다. 이 튜토리얼에서 사용된 입력 데이터 파일은 CSV 파일인 SalesJan2009.csv입니다.

public static class MapClass
            extends MapReduceBase
            implements Mapper<LongWritable, Text, Text, Text>
{
    static enum SalesCounters { MISSING, INVALID };
    public void map ( LongWritable key, Text value,
                 OutputCollector<Text, Text> output,
                 Reporter reporter) throws IOException
    {
        
        //Input string is split using ',' and stored in 'fields' array
        String fields[] = value.toString().split(",", -20);
        //Value at 4th index is country. It is stored in 'country' variable
        String country = fields[4];
        
        //Value at 8th index is sales data. It is stored in 'sales' variable
        String sales = fields[8];
      
        if (country.length() == 0) {
            reporter.incrCounter(SalesCounters.MISSING, 1);
        } else if (sales.startsWith("\"")) {
            reporter.incrCounter(SalesCounters.INVALID, 1);
        } else {
            output.collect(new Text(country), new Text(sales + ",1"));
        }
    }
}

위 코드 조각은 Hadoop MapReduce에서 카운터를 구현한 예시를 보여줍니다.

여기 판매 카운터 '를 사용하여 정의된 카운터입니다.열거 형'. 누락되거나 유효하지 않은 입력 레코드 수를 세는 데 사용됩니다.

코드 조각에서, 만약 '국가' 필드의 길이가 0이면 해당 값이 누락된 것이므로 그에 해당하는 카운터인 SalesCounters.MISSING이 증가합니다.

다음으로, 만약 '판매' 필드가 "로 시작하면 해당 레코드는 유효하지 않은 것으로 간주됩니다. 이는 SalesCounters.INVALID 카운터를 증가시켜 표시됩니다.

💡 API 참고: 위의 코드 조각은 원본을 사용합니다. org.apache.hadoop.mapred API, 여기서 MapReduceBase 밸리 Mapper 인터페이스 OutputCollector Reporter 별도로 표시됩니다. 현재 코드는 다음을 기준으로 작성되었습니다. org.apache.hadoop.mapreduce단일한 Context 수집기와 보고기를 대체하고 카운터가 증가합니다. context.getCounter(SalesCounters.MISSING).increment(1)카운터 개념은 두 경우 모두 동일합니다.

자주 묻는 질문

한쪽 데이터가 모든 노드의 메모리에 저장할 수 있을 만큼 작다면 셔플 과정을 완전히 생략하는 매퍼 방식을 선택하세요. 양쪽 데이터가 모두 크거나 정렬되지 않은 상태라면 추가적인 네트워크 비용을 감수하고 리듀서 방식을 선택하세요.

모델은 과거 작업 기록을 학습하여 실행 시간을 예측하고, 분할 크기와 리듀서 개수를 권장하며, 카운터 값에서 편차를 감지합니다. 또한 스필드 레코드 또는 실패한 작업 카운터가 해당 파이프라인의 정상 범위를 벗어나는 작업을 표시합니다.

Copilot은 그럴듯한 매퍼 및 리듀서 골격을 생성하지만, 하나의 클래스에서 기존 mapred 패키지와 최신 mapreduce 패키지를 자유롭게 혼합하여 사용하므로 컴파일되지 않습니다. 로직을 신뢰하기 전에 임포트 및 메서드 시그니처를 수정하십시오.

이는 작업이 시작되기 전에 더 작은 파일을 모든 노드로 전송하는 메커니즘입니다. 각 매퍼는 해당 복사본을 해시 맵에 로드하고 로컬에서 일치하는 항목을 찾는데, 이것이 매퍼 측 조인이 가능한 이유입니다.

이러한 정보는 작업이 완료되면 콘솔 요약에 출력되고, 작업 기록 및 리소스 관리자 웹 페이지에 표시되며, 작업 객체에서 프로그래밍 방식으로 읽을 수 있으므로 드라이버는 이러한 정보를 검증하고 잘못된 실행을 실패로 처리할 수 있습니다.

단일 Context 객체입니다. 이 객체는 OutputCollector와 Reporter가 이전에 분담했던 작업을 수행하므로, 출력은 map 메서드에 전달된 동일한 핸들을 통해 기록되고 카운터는 증가합니다.

대부분의 보도 업무에서는 그렇지 않습니다. A HiveQL 조인 몇 줄의 코드로 동일한 셔플 및 병합 패턴으로 컴파일됩니다. 병합 로직이 SQL 절에 맞지 않을 때는 직접 코드를 작성하는 것이 가치가 있습니다.

Hadoop은 이미 존재하는 출력 디렉터리에 쓰기를 거부합니다. 이는 완료된 결과가 덮어쓰이는 것을 방지하기 위한 조치입니다. 먼저 해당 디렉터리를 재귀적으로 삭제하거나, 다음 실행 시 다른 출력 경로를 지정하십시오.

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