리듀스 측 조인 (Reduce-Side Join)
리듀스 측 조인은 리파티셔닝 정렬-병합 조인(Repartitioned Sort-Merge Join)이라고도 불립니다. 이 방식은 다양한 데이터 소스의 레코드를 리듀서로 전송하여 결합하므로, 셔플(Shuffle) 단계에서 발생하는 네트워크 및 I/O 오버헤드가 크다는 특징이 있습니다.
핵심 개념 및 처리 흐름
- 데이터 소스 (Data Source): 관계형 데이터베이스의 테이블과 유사한 개념으로, 동일한 스키마를 가진 여러 파일로 구성될 수 있습니다.
- 태그 (Tag): 레코드가 어떤 데이터 소스에서 유래했는지 식별하기 위한 메타데이터(예: 파일 이름)입니다.
- 그룹 키 (Group Key): 조인 연산을 수행하기 위한 기준이 되는 키(Join Key)입니다.
처리 프로세스는 다음과 같이 진행됩니다.
- Map 단계: 각 데이터 소스에서 레코드를 읽어 그룹 키와 출처를 식별하는 태그로 패키징합니다.
- Shuffle 단계: 그룹 키를 기준으로 레코드를 정렬하고 동일한 키를 가진 레코드들을 하나의 그룹으로 묶습니다.
- Reduce 단계: 그룹화된 레코드들을 언패킹하여 원본 데이터를 복원한 뒤, 태그를 기준으로 교차 곱(Cartesian Product)을 수행합니다.
- Combine 단계: 내부 조인(Inner Join) 또는 외부 조인(Outer Join) 조건에 따라 교차 곱의 결과를 병합하여 최종 레코드를 생성합니다.
커스텀 패키징을 통한 구현
과거 Hadoop의 datajoin 패키지는 TaggedMapOutput, DataJoinMapperBase 등의 추상 클래스를 제공했으나 현재는 지원이 중단되었습니다. 따라서 최신 MapReduce API를 활용하여 태그와 데이터를 패키징하는 커스텀 Writable 클래스를 구현하는 것이 일반적입니다.
// 데이터와 출처 태그를 묶는 커스텀 Writable 클래스
public class TaggedRecord implements Writable {
private Text sourceTag;
private Text payload;
public TaggedRecord() {
this.sourceTag = new Text();
this.payload = new Text();
}
public TaggedRecord(String tag, String data) {
this.sourceTag = new Text(tag);
this.payload = new Text(data);
}
public Text getSourceTag() { return sourceTag; }
public Text getPayload() { return payload; }
@Override
public void write(DataOutput out) throws IOException {
sourceTag.write(out);
payload.write(out);
}
@Override
public void readFields(DataInput in) throws IOException {
sourceTag.readFields(in);
payload.readFields(in);
}
}
// 맵퍼 구현: 파일 이름을 태그로 활용하여 레코드 패키징
public class JoinMapper extends Mapper {
private Text joinKey = new Text();
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
String[] fields = value.toString().split(",");
if (fields.length > 1) {
joinKey.set(fields[0]); // 첫 번째 필드를 조인 키로 설정
String fileName = ((FileSplit) context.getInputSplit()).getPath().getName();
context.write(joinKey, new TaggedRecord(fileName, value.toString()));
}
}
}
맵 측 조인 (Map-Side Replicated Join)
조인에 필요한 모든 데이터를 맵퍼의 메모리 내에서 접근할 수 있다면, 셔플 단계를 완전히 생략하고 맵 단계에서만 조인을 수행할 수 있습니다. 이를 복제 조인(Replicated Join)이라고 하며, 일반적으로 하나의 테이블은 매우 크고 다른 하나의 테이블은 작아 메모리에 적재 가능할 때 사용됩니다.
Hadoop의 분산 캐시(Distributed Cache) 메커니즘을 사용하면 작은 차원 테이블을 클러스터의 모든 노드 로컬 디스크에 복제할 수 있습니다.
구현 예시: 사용자 정보와 거래 내역 조인
작은 테이블인 users.csv를 분산 캐시에 등록하고, 큰 테이블인 transactions.csv를 맵퍼에서 읽으며 메모리 내 해시맵과 조인하는 예시입니다.
public class ReplicatedJoinMapper extends Mapper {
private HashMap userLookup = new HashMap<>();
private Text outputValue = new Text();
@Override
protected void setup(Context context) throws IOException, InterruptedException {
// 분산 캐시에서 작은 테이블 파일 경로 가져오기
URI[] cachedFiles = context.getCacheFiles();
if (cachedFiles != null && cachedFiles.length > 0) {
Path smallTablePath = new Path(cachedFiles[0]);
// 로컬 파일 시스템에서 작은 테이블을 읽어 메모리에 적재
try (BufferedReader reader = new BufferedReader(new FileReader(smallTablePath.toString()))) {
String line;
while ((line = reader.readLine()) != null) {
String[] parts = line.split(",", 2);
if (parts.length == 2) {
userLookup.put(parts[0], parts[1]);
}
}
}
}
}
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
String[] fields = value.toString().split(",");
String userId = fields[0];
String matchedUserInfo = userLookup.get(userId);
// 메모리 내 해시맵과 조인 수행
if (matchedUserInfo != null) {
outputValue.set(value.toString() + "\t" + matchedUserInfo);
context.write(new Text(userId), outputValue);
}
}
}
// 드라이버 클래스에서 캐시 파일 등록
// job.addCacheFile(new URI("hdfs://namenode/path/to/users.csv#users"));
세미 조인 (Semi-Join)
두 데이터 소스 모두 규모가 커서 맵 측 조인을 적용할 수 없는 경우, 세미 조인 기법을 사용하여 네트워크 전송량을 획기적으로 줄일 수 있습니다. 세미 조인은 리듀스로 데이터를 전송하기 전에 맵 단계에서 불필요한 레코드를 사전에 필터링하는 방식입니다.
일반적으로 블룸 필터(Bloom Filter)를 활용하여 조인 키의 존재 여부를 빠르게 확인합니다. 작은 테이블의 조인 키들을 바탕으로 블룸 필터를 생성한 뒤, 이를 분산 캐시를 통해 모든 맵퍼에 배포합니다. 맵퍼는 큰 테이블을 스캔하며 블룸 필터를 통과한 레코드만 리듀스로 전달하고, 최종적인 정확한 조인 연산은 리듀스 측 조인 방식을 통해 수행됩니다. 이를 통해 셔플 단계로 전달되는 데이터의 볼륨을 최소화하여 전체 잡의 성능을 향상시킬 수 있습니다.