S3 에 있는 파일을 HDFS 로 내려 가공한 뒤, 원본을 지우고 결과를 다시 S3 로 올리는 처리를 자바 애플리케이션 안에서 하는 경우다. distcp 는 대량 일괄 전송에 맞고, 파일 단위로 조건을 보며 옮기는 처리에는 org.apache.hadoop.fs API 가 맞다.
핵심은 하나다. AWS SDK 를 직접 쓸 필요가 없다. s3a:// 스킴을 쓰면 S3 도 HDFS 와 같은 FileSystem 추상화로 다뤄지므로, 목록 조회 · 복사 · 삭제를 같은 코드로 처리할 수 있다.
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>${hadoop.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-aws</artifactId>
<version>${hadoop.version}</version>
</dependency>
hadoop-aws 버전은 hadoop-client 와 반드시 같게 맞춘다. 다르면 런타임에 NoSuchMethodError 가 난다. AWS SDK 는 hadoop-aws 가 끌고 오는 것을 그대로 쓰고 직접 버전을 지정하지 않는다.
fs.s3a.aws.credentials.provider 로 자격 증명 공급자를 지정하면 키를 코드나 설정 파일에 두지 않아도 된다.
| 공급자 | 읽는 곳 |
|---|---|
EnvironmentVariableCredentialsProvider |
AWS_ACCESS_KEY_ID · AWS_SECRET_ACCESS_KEY 환경 변수 |
InstanceProfileCredentialsProvider |
EC2 인스턴스 프로파일(IAM 역할) |
WebIdentityTokenCredentialsProvider |
EKS 의 IRSA |
SimpleAWSCredentialsProvider |
fs.s3a.access.key · fs.s3a.secret.key |
EC2 · EKS 에서 돌린다면 IAM 역할을 붙이는 것이 가장 낫다. 키를 발급하지 않으므로 유출과 회전 관리가 사라진다. 환경 변수로 준다면 애플리케이션을 띄우는 셸이나 systemd 유닛(Environment=), 컨테이너의 env 에 넣는다. .bashrc 에 넣어 두면 배치 실행 환경에서 읽히지 않는 경우가 많다.
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.FileUtil;
import org.apache.hadoop.fs.Path;
import java.net.URI;
public class S3HdfsMover {
private static Configuration conf() {
Configuration c = new Configuration();
// core-site.xml 이 클래스패스에 있으면 대부분 생략 가능
c.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem");
c.set("fs.s3a.aws.credentials.provider",
"com.amazonaws.auth.EnvironmentVariableCredentialsProvider");
c.set("fs.s3a.endpoint", "s3.ap-northeast-2.amazonaws.com");
return c;
}
public static void main(String[] args) throws Exception {
Configuration conf = conf();
Path src = new Path("s3a://my-bucket/raw/2023-05-26/");
Path work = new Path("hdfs://nameservice1/staging/2023-05-26/");
Path dst = new Path("s3a://my-bucket/curated/2023-05-26/");
FileSystem s3 = FileSystem.get(URI.create(src.toString()), conf);
FileSystem hdfs = FileSystem.get(URI.create(work.toString()), conf);
// 1. 목록 조회 — 파일만 고른다
for (FileStatus st : s3.listStatus(src)) {
if (st.isDirectory()) {
continue;
}
Path in = st.getPath();
Path staged = new Path(work, in.getName());
// 2. S3 → HDFS (원본 유지)
FileUtil.copy(s3, in, hdfs, staged, false, true, conf);
// 3. 가공 (생략)
// 4. HDFS → S3
Path out = new Path(dst, in.getName());
FileUtil.copy(hdfs, staged, s3, out, false, true, conf);
// 5. 확인 후 원본 삭제
if (s3.exists(out) && s3.getFileStatus(out).getLen() > 0) {
s3.delete(in, false);
}
}
hdfs.close();
s3.close();
}
}
FileUtil.copy(srcFS, src, dstFS, dst, deleteSource, overwrite, conf) 의 인자 순서를 기억해 둔다. 다섯 번째가 원본 삭제 여부, 여섯 번째가 덮어쓰기다. 삭제는 별도 단계로 두는 편이 안전하다. deleteSource=true 로 한 번에 옮기면 대상 쓰기가 실패했는데 원본이 사라지는 경우가 생긴다.
디렉터리를 통째로 옮길 때도 같은 호출로 재귀 복사된다. 다만 파일이 많으면 한 건씩 순차 복사라 느리므로, 그 규모에서는 distcp 가 낫다.
s3a 가 접두사를 디렉터리처럼 보여 줄 뿐이다. 빈 디렉터리를 만들거나 rename 을 기대하면 예상과 다르게 동작한다. rename 은 복사 후 삭제로 구현되므로 큰 파일에서는 비용이 크고 원자적이지 않다.delete(path, true) 는 재귀 삭제다. 접두사를 한 글자 잘못 적으면 되돌릴 수 없다. 실행 전에 listStatus 로 대상을 출력해 눈으로 확인하는 단계를 코드에 넣는다.getFileChecksum 과 S3 의 ETag 는 알고리즘이 달라 값을 직접 비교할 수 없다. 크기 비교 또는 애플리케이션이 계산한 해시를 메타데이터로 남겨 비교한다.listStatus 를 걸면 페이지 요청이 반복돼 느리다. 날짜 등으로 접두사를 잘게 나눈다.가공까지 Spark 로 한다면 FileSystem API 로 내려받을 필요가 없다. 읽고 쓰고 지우는 것을 모두 경로로 처리한다.
val df = spark.read.parquet("s3a://my-bucket/raw/2023-05-26/")
val out = transform(df)
out.write.mode("overwrite").parquet("s3a://my-bucket/curated/2023-05-26/")
// 삭제만 FileSystem API 로
val fs = org.apache.hadoop.fs.FileSystem.get(
java.net.URI.create("s3a://my-bucket/"), spark.sparkContext.hadoopConfiguration)
fs.delete(new org.apache.hadoop.fs.Path("s3a://my-bucket/raw/2023-05-26/"), true)
이 방식에서는 중간 HDFS 스테이징이 필요 없다. 쓰기 커밋 방식(fs.s3a.committer.name)을 매직 커미터로 두면 이름 변경 비용도 피할 수 있다.