Amazon S3
Amazon S3
Amazon Simple Storage Service (Amazon S3)는 다양한 사용 사례를 위한 클라우드 객체 스토리지를 제공해요. Flink와 함께 S3를 데이터 읽기와 쓰기에 사용할 수 있고, 스트리밍 state backends와 함께 사용할 수도 있어요.
출처: 문서
본문
다음 형식으로 경로를 지정해 S3 객체를 일반 파일처럼 사용할 수 있어요:
s3://<your-bucket>/<endpoint>
끝점(endpoint)은 단일 파일이나 디렉터리가 될 수 있어요. 예를 들어:
// Read from S3 bucket
FileSource<String> fileSource = FileSource.forRecordStreamFormat(
new TextLineInputFormat(), new Path("s3://<your-bucket>/<endpoint>")
).build();
env.fromSource(
fileSource,
WatermarkStrategy.noWatermarks(),
"s3-input"
);
// Write to S3 bucket
stream.sinkTo(
FileSink.forRowFormat(
new Path("s3://<your-bucket>/<endpoint>"), new SimpleStringEncoder<>()
).build()
);
// Use S3 as checkpoint storage
Configuration config = new Configuration();
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "s3://<your-bucket>/<endpoint>");
env.configure(config);
이 예시들은 완전한 것이 아닙니다. Flink가 FileSystem URI를 기대하는 모든 곳(달리 명시되지 않는 한), 예를 들어 고가용성(HA) 설정이나 EmbeddedRocksDBStateBackend에서도 S3를 사용할 수 있어요.
S3 파일시스템 구현 (S3 FileSystem Implementations)
Flink는 세 가지 독립적인 S3 파일시스템 구현을 제공해요:
| 구현 | Checkpointing | FileSink | 참고(Notes) |
|---|---|---|---|
Native S3 (flink-s3-fs-native) |
✓ | ✓ | Flink 2.3에서 실험적(Experimental). AWS SDK v2 기반, Hadoop 의존성 없음. |
Presto S3 (flink-s3-fs-presto) |
✓ | x | 체크포인팅에 대해 프로덕션 검증됨. |
Hadoop S3 (flink-s3-fs-hadoop) |
✓ | ✓ | 성숙함. FileSink를 위한 RecoverableWriter를 제공하는 유일한 안정적 구현. |
이전에는 사용자가 Presto(체크포인팅 처리량에 권장)와 Hadoop(RecoverableWriter를 가진 유일한 구현으로, FileSink에 필요) 중에서 선택해야 했어요. Native S3 구현은 두 기능을 단일 플러그인으로 통합하며, 벤치마크 결과 Presto 구현보다 체크포인팅 처리량이 크게 개선되는 것으로 나타나요.
세 구현 모두 자체 완결적(self-contained)이며 의존성 풋프린트가 없으므로, 사용하기 위해 Hadoop을 클래스패스에 추가할 필요가 없어요.
공통 구성 (Common Configuration)
접근 자격 증명 구성 (Configure Access Credentials)
S3 파일시스템 구현을 설정한 후, Flink가 S3 버킷에 접근할 수 있도록 해야 해요. 다음 세 가지 접근 방식은 독립적인 대안이에요 — 환경에 맞는 것을 선택하세요:
IAM (Identity and Access Management) (권장)
AWS에서 자격 증명을 설정하는 권장 방법은 Identity and Access Management (IAM)을 통한 것이에요. IAM 기능을 사용해 Flink 인스턴스에 S3 버킷 접근에 필요한 자격 증명을 안전하게 부여할 수 있어요. 이 작업 방법에 대한 자세한 내용은 이 문서의 범위를 벗어나요. AWS 사용자 가이드를 참조하세요. 찾고 있는 것은 IAM Roles예요.
이것을 올바르게 설정하면 AWS 내에서 S3에 대한 접근을 관리할 수 있고 Flink에 접근 키를 배포할 필요가 없어요.
위임 토큰 (Delegation Tokens)
위임 토큰은 시간 제한이 있고 자동으로 협상되는 자격 증명을 제공해요. JobManager는 장기 수명의 자격 증명(액세스 키와 시크릿 키)을 사용해 AWS STS를 호출하고, 단기 수명의 세션 토큰을 얻어 자동으로 TaskManager에 배포해요.
각 S3 구현은 전용 구성 접두사를 가진 고유한 위임 토큰 제공자(provider)를 가져요. 사용 중인 구현의 해당 접두사 아래에 access-key, secret-key, region을 설정해야 해요:
# For Native S3 implementation
security.delegation.token.provider.s3-native.access-key: your-access-key
security.delegation.token.provider.s3-native.secret-key: your-secret-key
security.delegation.token.provider.s3-native.region: us-east-1
# For Hadoop implementation
security.delegation.token.provider.s3-hadoop.access-key: your-access-key
security.delegation.token.provider.s3-hadoop.secret-key: your-secret-key
security.delegation.token.provider.s3-hadoop.region: us-east-1
# For Presto implementation
security.delegation.token.provider.s3-presto.access-key: your-access-key
security.delegation.token.provider.s3-presto.secret-key: your-secret-key
security.delegation.token.provider.s3-presto.region: us-east-1
위임 토큰이 발급되려면 세 값(access-key, secret-key, region)이 모두 설정되어야 해요. DynamicTemporaryAWSCredentialsProvider가 각 구현의 자격 증명 제공자 체인에 자동으로 포함되므로, TaskManager는 추가 구성 없이 배포된 토큰을 소비해요.
접근 키 (Access Keys)
S3 접근은 접근 키와 시크릿 키 쌍을 통해 부여될 수 있어요. 접근 키 자체가 본질적으로 안전하지 않은 것은 아니지만, 정적 자격 증명을 관리/배포할 필요를 피하므로 IAM 역할이 선호돼요. 더 많은 맥락은 IAM 역할 소개를 참조하세요.
Flink의 설정 파일에서 s3.access-key와 s3.secret-key를 모두 구성해야 해요:
s3.access-key: your-access-key
s3.secret-key: your-secret-key
비-S3 끝점 구성 (Configure Non-S3 Endpoint)
S3 파일시스템은 S3 호환 객체 스토어(compliant object stores)도 지원해요. 그러려면 Flink 설정 파일에서 끝점을 구성해요:
s3.endpoint: your-endpoint-hostname
경로 스타일 접근 구성 (Configure Path Style Access)
일부 S3 호환 객체 스토어는 기본적으로 가상 호스트 스타일 어드레싱(virtual host style addressing)이 활성화되지 않을 수 있어요. 이런 경우 Flink 설정 파일에서 경로 스타일 접근을 활성화하는 속성을 제공해야 해요:
s3.path-style-access: true
레거시 구성 키
s3.path.style.access는 하위 호환성을 위한 폴백으로 여전히 지원돼요.
구현 세부 사항 (Implementation Details)
Native S3 파일시스템 (실험적)
실험적: Native S3 FileSystem은 Flink 2.3에서 실험적이에요. 기능적으로 완전하며 벤치마크에서 강력한 성능을 입증했어요.
Native S3 파일시스템은 AWS SDK v2 위에 구축된 순수 자바 구현으로, Hadoop에 대한 의존성을 완전히 제거했어요. s3:// 및 s3a:// 스킴으로 등록돼요. Presto와 Hadoop 구현을 대체(drop-in replacement)하는 것으로, 체크포인팅, FileSink(RecoverableWriter를 통해), 서버 측 암호화(SSE-S3, SSE-KMS), IAM 역할 위임을 통한 크로스 계정 접근, 엔트로피 주입(entropy injection), S3TransferManager를 통한 일괄 복사(bulk copy)를 지원해요.
설치 (Setup)
Native S3 파일시스템을 사용하려면 opt 디렉터리의 JAR 파일을 plugins 디렉터리로 복사해요:
mkdir -p ./plugins/s3-fs-native
cp ./opt/flink-s3-fs-native-2.3.0.jar ./plugins/s3-fs-native/
구성 (Configuration)
공통 구성 옵션(s3.access-key, s3.secret-key, s3.endpoint, s3.path-style-access)에 더해, Native S3 파일시스템은 다음 옵션을 지원해요:
# Server-side encryption
s3.sse.type: sse-s3 # or sse-kms, aws:kms, AES256, none (default)
s3.sse.kms.key-id: arn:aws:kms:region:account:key/id # Required for SSE-KMS
# IAM role assumption for cross-account access
s3.assume-role.arn: arn:aws:iam::account:role/RoleName
s3.assume-role.external-id: external-id-if-required
s3.assume-role.session-name: flink-s3-session
s3.assume-role.session-duration: 3600
# Performance tuning
s3.upload.min.part.size: 5242880 # 5 MB default
s3.upload.max.concurrent.uploads: 4 # Based on CPU cores
s3.read.buffer.size: 262144 # 256 KB default
s3.async.enabled: true # Async read/write operations
s3.bulk-copy.enabled: true # Bulk copy via S3TransferManager
s3.bulk-copy.max-concurrent: 16 # Max concurrent copy ops
fs.s3.aws.credentials.provider가 설정되지 않으면, Native S3 파일시스템은 다음과 같은 순서로 자격 증명 체인을 자동으로 구성해요: 위임 토큰, 정적 자격 증명(s3.access-key와 s3.secret-key가 구성된 경우), 그리고 AWS SDK v2 DefaultCredentialsProvider(환경 변수, 인스턴스 프로파일 등). 커스텀 제공자 체인이 필요한 경우에만 이 옵션을 설정하면 돼요.
Presto S3 파일시스템
EMR에서 Flink를 실행한다면 이걸 수동으로 구성할 필요가 없어요.
Presto S3 파일시스템은 Presto 프로젝트의 코드를 기반으로 해요. s3:// 및 s3p:// 스킴으로 등록돼요. S3로의 체크포인팅에 프로덕션에서 검증된 선택이에요. FileSink(createRecoverableWriter는 UnsupportedOperationException을 던짐)는 지원하지 않아요.
설치 (Setup)
Presto S3 파일시스템을 사용하려면 opt 디렉터리의 JAR 파일을 plugins 디렉터리로 복사해요:
mkdir -p ./plugins/s3-fs-presto
cp ./opt/flink-s3-fs-presto-2.3.0.jar ./plugins/s3-fs-presto/
구성 (Configuration)
공통 구성 옵션이 적용돼요. 추가로 Presto 전용 키는 Presto 파일 시스템 구성을 통해 지원돼요.
Hadoop S3 파일시스템
Hadoop S3 파일시스템은 Hadoop 프로젝트의 코드를 기반으로 해요. s3:// 및 s3a:// 스킴으로 등록돼요. FileSink(RecoverableWriter를 통해)를 지원하는 유일한 안정적 구현이에요.
설치 (Setup)
Hadoop S3 파일시스템을 사용하려면 opt 디렉터리의 JAR 파일을 plugins 디렉터리로 복사해요:
mkdir -p ./plugins/s3-fs-hadoop
cp ./opt/flink-s3-fs-hadoop-2.3.0.jar ./plugins/s3-fs-hadoop/
구성 (Configuration)
공통 구성 옵션이 적용돼요. 추가로 Hadoop의 s3a 구성 키가 지원돼요. Hadoop 구성 키는 자동으로 번역돼요 — 예를 들어 fs.s3a.connection.maximum은 s3.connection.maximum이 돼요.
여러 S3 구현 사용하기 (Using Multiple S3 Implementations)
세 S3 구현 모두 s3:// 스킴의 핸들러로 등록돼요. 추가로 각 구현은 대체 스킴을 지원해요:
| 구현 | 스킴(Schemes) |
|---|---|
| Native S3 | s3://, s3a:// |
| Presto | s3://, s3p:// |
| Hadoop | s3://, s3a:// |
여러 S3 플러그인 JAR을 동시에 로드해도 안전해요 — 우선순위 메커니즘이 각 스킴을 처리하는 팩토리가 하나만 있도록 보장해요. Native S3 구현은 가장 낮은 우선순위(-1 vs 기본 0)를 가지므로, 다른 구현이 존재하면 겹치는 스킴(예: s3:// 및 s3a://)에 대해 그것이 우선해요. fs.<scheme>.priority.<factory> 구성 옵션으로 팩토리 우선순위를 재정의할 수 있어요.
서로 다른 URI 스킴을 활용해 여러 S3 구현을 동시에 사용할 수 있어요. 예를 들어, 작업이 FileSystem 싱크는 Hadoop으로, 체크포인팅은 Presto로 사용한다면:
- 싱크에 s3a:// 스킴 사용 (Hadoop)
- 체크포인팅에 s3p:// 스킴 사용 (Presto)
Native S3 구현은 새 URI 스킴을 도입하지 않아요. 기존 s3:// 및 s3a:// 스킴을 지원해요. Native S3와 Hadoop 구현이 동일한 스킴으로 등록되므로, Flink는 각 스킴을 처리할 팩토리를 선택하기 위해 우선순위 기반 메커니즘을 사용해요. 기본적으로 Native S3는 가장 낮은 우선순위를 가지며, 동일한 스킴에 다른 구현이 존재하면 선택되지 않아요.
Native S3 구현을 사용하려면
plugins디렉터리에flink-s3-fs-native플러그인 JAR만 두거나, 다른 구현이plugins에 있는 동안fs.<scheme>.priority.<factory>구성을 사용해 그 우선순위를 올리세요.
고급 기능 (Advanced Features)
엔트로피 주입 (Entropy Injection)
모든 S3 파일 시스템은 엔트로피 주입을 지원해요. 키(key)의 시작 부분 근처에 무작위 문자를 추가해 AWS S3 버킷의 확장성을 개선하는 기술이에요.
엔트로피 주입이 활성화되면, 경로의 구성된 부분 문자열이 무작위 문자로 대체돼요. 예를 들어 경로 s3://my-bucket/_entropy_/checkpoints/dashboard-job/는 s3://my-bucket/gf36ikvg/checkpoints/dashboard-job/ 같은 것으로 대체돼요.
이것은 파일 생성이 엔트로피 주입 옵션을 전달할 때만 발생해요!
그렇지 않으면 파일 경로가 엔트로피 키 부분 문자열을 완전히 제거해요. 자세한 내용은 FileSystem.create(Path, WriteOption)를 참조하세요.
Flink 런타임은 현재 엔트로피 주입 옵션을 체크포인트 데이터 파일에만 전달해요. 체크포인트 메타데이터와 외부 URI를 포함한 다른 모든 파일은 체크포인트 URI를 예측 가능하게 유지하기 위해 엔트로피를 주입하지 않아요.
엔트로피 주입을 활성화하려면 *엔트로피 키(entropy key)*와 엔트로피 길이(entropy length) 파라미터를 구성해요.
s3.entropy.key: _entropy_
s3.entropy.length: 4 (default)
s3.entropy.key는 무작위 문자로 대체되는 경로의 문자열을 정의해요. 엔트로피 키를 포함하지 않는 경로는 변경되지 않아요.
파일 시스템 연산이 "inject entropy" 쓰기 옵션을 전달하지 않으면, 엔트로피 키 부분 문자열은 단순히 제거돼요.
s3.entropy.length는 엔트로피에 사용되는 무작위 영숫자 문자 수를 정의해요.
s5cmd
지원: Presto S3 파일시스템, Hadoop S3 파일시스템
flink-s3-fs-hadoop과 flink-s3-fs-presto 모두 더 빠른 파일 업로드/다운로드를 위해 s5cmd 툴을 사용하도록 구성할 수 있어요.
벤치마크 결과에 따르면 s5cmd는 CPU 효율이 2배 이상 높을 수 있어요.
즉 같은 파일 집합을 업로드/다운로드하는 데 CPU를 절반만 사용하거나, 같은 양의 CPU로 두 배 빠르게 수행한다는 뜻이에요.
이 기능을 사용하려면 s5cmd 바이너리가 존재하고 Flink의 task manager가 접근할 수 있어야 해요. 예를 들어 사용하는 docker 이미지에 내장하는 식으로요.
둘째로 s5cmd의 경로를 다음과 같이 구성해야 해요:
s3.s5cmd.path: /path/to/the/s5cmd
구성 (Configuration)
남은 구성 옵션(기본값은 아래)은 다음과 같아요:
# Extra arguments that will be passed directly to the s5cmd call. Please refer to the s5cmd's official documentation.
s3.s5cmd.args: -r 0
# Maximum size of files that will be uploaded via a single s5cmd call.
s3.s5cmd.batch.max-size: 1024mb
# Maximum number of files that will be uploaded via a single s5cmd call.
s3.s5cmd.batch.max-files: 100
s3.s5cmd.batch.max-size과 s3.s5cmd.batch.max-files 모두 task manager에 과부하가 걸리지 않도록 s5cmd 바이너리의 리소스 사용을 제어해요.
먼저 s5cmd 없이 Flink가 동작하는지 구성하고 검증한 다음 이 기능을 활성화하는 것이 권장돼요.
자격 증명 (Credentials)
접근 키를 사용한다면, 그것들이 s5cmd에 전달돼요.
그 외에 s5cmd는 자체적으로 독립적인 자격 증명 사용 방식을 가져요.
제한 사항 (Limitations)
현재 flink-s3-fs-hadoop과 flink-s3-fs-presto는 복구 중 상태 파일을 S3에서 다운로드하고 RocksDB를 사용할 때만 s5cmd를 사용해요.
flink-s3-fs-native는 s3.bulk-copy.enabled(기본값: true)로 활성화하면 일괄 복사 연산에 S3TransferManager를, s3.async.enabled(기본값: true)로 비동기 읽기/쓰기에 S3TransferManager를 사용해 유사한 성능 이점을 제공해요.