Google Cloud Storage
Google Cloud Storage
Google Cloud Storage(GCS)는 다양한 사용 사례를 위한 클라우드 저장소를 제공해요. 스트리밍 상태 백엔드와 함께 FileSystemCheckpointStorage를 사용할 때 데이터 읽기·쓰기와 체크포인트 저장소로 사용할 수 있어요.
출처: 문서
본문
다음 형식으로 경로를 지정하면 GCS 객체를 일반 파일처럼 사용할 수 있어요.
gs://<your-bucket>/<endpoint>
엔드포인트는 단일 파일일 수도 있고 디렉터리일 수도 있어요. 예를 들어:
// Read from GCS bucket
FileSource<String> fileSource = FileSource.forRecordStreamFormat(
new TextLineInputFormat(),
new Path("gs://<bucket>/<endpoint>")
).build();
env.fromSource(
fileSource,
WatermarkStrategy.noWatermarks(),
"gcs-input"
);
// Write to GCS bucket
stream.sinkTo(
FileSink.forRowFormat(
new Path("gs://<bucket>/<endpoint>"),
new SimpleStringEncoder<>()
).build()
);
// Use GCS as checkpoint storage
Configuration config = new Configuration();
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY, "gs://<bucket>/<endpoint>");
env.configure(config);
이 예제는 완전하지 않으며, Flink가 FileSystem URI를 기대하는 모든 곳(고가용성 설정이나 EmbeddedRocksDBStateBackend 포함)에서 GCS를 사용할 수 있어요.
GCS 파일 시스템 플러그인 (GCS File System plugin)
Flink는 GCS에 쓰기 위한 flink-gs-fs-hadoop 파일 시스템을 제공해요. 이 구현은 자체 포함이라 의존성 풋프린트가 없어 사용하기 위해 Hadoop을 클래스패스에 추가할 필요가 없어요.
flink-gs-fs-hadoop은 gs:// 스킴의 URI에 대한 FileSystem 래퍼를 등록해요. GCS 접근에는 Google의 gcs-connector Hadoop 라이브러리를 사용하고, RecoverableWriter 지원을 위해 Google의 google-cloud-storage 라이브러리를 사용해요.
이 파일 시스템은 FileSystem 커넥터와 함께 사용할 수 있어요.
flink-gs-fs-hadoop을 사용하려면 Flink를 시작하기 전에 JAR 파일을 opt 디렉터리에서 Flink 배포의 plugins 디렉터리로 복사해요:
mkdir ./plugins/gs-fs-hadoop
cp ./opt/flink-gs-fs-hadoop-2.3.0.jar ./plugins/gs-fs-hadoop/
구성 (Configuration)
기본 Hadoop 파일 시스템은 Hadoop 구성 키를 Flink 구성 파일에 추가해 구성할 수 있어요. 예를 들어 gcs-connector에 fs.gs.http.connect-timeout 구성 키가 있다면, Flink 구성 파일에서 gs.http.connect-timeout: xyz로 설정하면 돼요. Flink는 이를 내부적으로 fs.gs.http.connect-timeout으로 번역해요.
또한 env.hadoop.conf.dir Flink 옵션이나 HADOOP_CONF_DIR 환경 변수를 통해 Hadoop 구성 디렉터리를 Flink에 알리면, gcs-connector 옵션을 Hadoop core-site.xml 구성 파일에 직접 설정할 수도 있어요.
flink-gs-fs-hadoop은 Flink 구성 파일에서 다음 옵션을 설정해도 구성할 수 있어요:
| 키 (Key) | 설명 (Description) |
|---|---|
| gs.writer.temporary.bucket.name | RecoverableWriter로 진행 중인 쓰기를 위한 임시 blob을 담을 버킷을 고르는 프로퍼티예요. 설정하지 않으면 임시 blob은 최종 파일과 같은 버킷에 쓰여져요. 어느 경우든 임시 blob은 .inprogress/ 접두사로 쓰여져요. 체크/세이브포인트에서 복원할 때 발생할 수 있는 고아 blob을 정리하는 메커니즘을 제공하기 위해 별도 버킷을 골라 TTL을 지정하는 것을 권장해요. 임시 blob용 별도 버킷에 TTL을 사용하면 TTL 간격이 지난 후 체크/세이브포인트에서 작업을 재시작하려 할 때 실패할 수 있어요. |
| gs.writer.chunk.size | RecoverableWriter를 통한 쓰기의 청크 크기를 설정하는 프로퍼티예요. 설정하지 않으면 Google이 결정한 기본 청크 크기가 사용돼요. |
| gs.filesink.entropy.enabled | GCS의 핫스팟 문제로 인한 성능을 개선하는 프로퍼티예요. filesink gcs 경로에 엔트로피 주입을 활성화할지 정의해요. 활성화하면 임시 객체 ID 형태의 엔트로피가 임시 객체의 gcs 경로 시작 부분에 주입돼요. 최종 객체 경로는 바뀌지 않아요. |
| gs.http.connect-timeout | java-storage 클라이언트의 연결 타임아웃을 설정하는 프로퍼티예요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.http.read-timeout | java-storage 클라이언트를 통해 확립된 연결에서 콘텐츠 읽기 타임아웃을 설정하는 프로퍼티예요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.retry.max-attempt | 수행할 최대 재시도 횟수를 정의하는 프로퍼티예요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.retry.init-rpc-timeout | 초기 RPC의 타임아웃을 설정하는 프로퍼티예요. 이후 호출은 gs.retry.rpc-timeout-multiplier에 따라 조정된 이 값을 사용해요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.retry.rpc-timeout-multiplier | RPC 타임아웃의 변화를 제어하는 프로퍼티예요. 이전 호출의 타임아웃에 RpcTimeoutMultiplier를 곱해 다음 호출의 타임아웃을 계산해요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.retry.max-rpc-timeout | RPC 타임아웃 값에 상한을 두는 프로퍼티예요. 최대 rpc 타임아웃이 이 값보다 높게 RPC 타임아웃을 늘릴 수 없어요. 설정하지 않으면 GCS 기본값이 사용돼요. |
| gs.retry.total-timeout | 재시도를 시도할 수 있는 총 기간을 변경하는 프로퍼티예요. 설정하지 않으면 GCS 기본값이 사용돼요. |
GCS 접근 인증 (Authentication to access GCS)
GCS의 대부분의 작업은 인증이 필요해요. 인증 자격 증명을 제공하려면 다음 중 하나를 해요:
- JobManager와 TaskManager가 실행되는 곳에서 여기에 설명된 대로
GOOGLE_APPLICATION_CREDENTIALS환경 변수를 JSON 자격 증명 파일 경로로 설정해요. 이 방법이 권장돼요. core-site.xml의google.cloud.auth.service.account.json.keyfile프로퍼티를 JSON 자격 증명 파일 경로로 설정해요(Hadoop 구성 디렉터리를 위에서 설명한 대로 Flink에 지정했는지 확인해요):
<configuration>
<property>
<name>google.cloud.auth.service.account.json.keyfile</name>
<value>PATH TO GOOGLE AUTHENTICATION JSON FILE</value>
</property>
</configuration>
flink-gs-fs-hadoop이 이 두 방법 중 하나로 자격 증명을 사용하려면 인증에 service account 사용이 활성화되어야 해요. 기본적으로 활성화되어 있지만, core-site.xml에서 다음을 설정해 비활성화할 수 있어요:
<configuration>
<property>
<name>google.cloud.auth.service.account.enable</name>
<value>false</value>
</property>
</configuration>
gcs-connector는 위에서 설명한 google.cloud.auth.service.account.json.keyfile 옵션 외에도 자격 증명을 제공하는 추가 옵션을 지원해요. 하지만 다른 옵션을 사용하면 제공된 자격 증명이 RecoverableWriter를 지원하는 google-cloud-storage 라이브러리에서 사용되지 않아, Flink의 복구 가능한 쓰기(recoverable-write) 작업이 실패할 것으로 예상돼요. 따라서 google.cloud.auth.service.account.json.keyfile 외의 gcs-connector 인증 자격 증명 옵션 사용은 권장하지 않아요.