본문으로 건너뛰기

Firebase-BigQuery-GCS 데이터 마이그레이션 기술 방안

최근 업데이트 2026. 10. 07.

최근 업데이트 시간: 2022-04-13

소개​

Firebase는 통합 도구로, 여러 애플리케이션과 손쉽게 통합할 수 있습니다. 공식 문서는 다음과 같습니다

https://firebase.google.com/docs

플랫폼별 Firebase 문서 보기

제품별 Firebase 문서 보기

Firebase 프로젝트의 데이터는 여러 곳으로 내보낼 수 있습니다. Google 산하 서비스이기 때문에 많은 사용자가 데이터를 클라우드 스토리지 웨어하우스, 즉 BigQuery로 내보내고, BigQuery의 강력한 sql 기능을 활용하여 데이터를 분석합니다. 따라서 일반적으로 Firebase + BigQuery가 대부분 개발자의 선택입니다.

Firebase 프로젝트의 데이터를 BigQuery로 내보내는 작업은 Firebase 사용자 측에서 수행해야 합니다. Firebase 공식 문서의 이 글을 참고할 수 있습니다.

https://firebase.google.com/docs/projects/bigquery-export

이 문서는 Firebase 프로젝트의 데이터를 BigQuery로 내보낸 후 BigQuery의 sql 조회 기능으로 데이터를 조회·필터링하고, 조회 결과 데이터를 클라우드 스토리지(즉 Google Cloud Storage)로 옮긴 다음, 클라우드 스토리지의 데이터 파일을 로컬 머신으로 다운로드하여 파일 데이터를 파싱·처리한 후 AE 플랫폼으로 전송하는 경우에만 적용됩니다.

전체 데이터 흐름은 아래 그림과 같습니다

원본 이미지 보기

1. 고객이 제공해야 하는 설정​

  1. Google BigQuery와 Google Cloud Storage에 대한 읽기/쓰기 권한이 있는 Service Account의 JSON 파일. 다음은 설정 예시입니다
{
"type": "service_account",
"project_id": "<your_project_id>",
"private_key_id": "<your_private_key_id>",
"private_key": "<your_private_key>",
"client_email": "<your_client_email>",
"client_id": "<your_client_id>",
"auth_uri": "https://accounts.google.com/o/oauth2/auth",
"token_uri": "https://oauth2.googleapis.com/token",
"auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs",
"client_x509_cert_url": "<your_client_x509_cert_url>"
}
  1. Google Cloud Storage의 bucketName(Google Cloud Storage와 Google BigQuery는 같은 region에 있어야 합니다). 다음은 설정 예시입니다
bucketName: "your-bucket-name"
  1. Google BigQuery의 프로젝트 이름. 특정 프로젝트가 BigQuery에 저장된 이름을 나타냅니다. 다음은 설정 예시입니다
projectId: "your-project-id"
  1. Google BigQuery의 특정 프로젝트에 속한 데이터 세트. 데이터 세트는 프로젝트 내의 구체적인 저장 위치를 나타내며, 데이터 세트 아래에는 일반적으로 데이터를 저장하는 테이블이 있습니다. 다음은 설정 예시입니다
dataset: "analytics_123456789"
  1. 데이터 테이블의 와일드카드. 한 프로젝트의 모든 데이터가 데이터 세트에 저장된 모든 테이블 이름에 같은 접두사가 있다면 이 접두사도 설정할 수 있습니다(이 접두사는 필수가 아님).

2. 샘플 프로젝트​

샘플 프로젝트는 현재 공개 다운로드 주소가 없습니다. ThinkingAI 고객 성공 매니저에게 문의하여 받으십시오.

2.1 샘플 프로젝트 구조​

2.2 코드 예시​

코드는 예시와 아이디어만 제공하므로 그대로 사용하지 마십시오!

이 코드 예시는 BigQuery의 sdk로 프로그램을 작성하는 방법만 보여 줍니다. 구체적인 비즈니스 로직은 실제 비즈니스에 맞게 수정해야 합니다!

2.2.1 BigQuery에서 데이터를 조회하고 결과를 클라우드 스토리지에 저장​

그 전에 Firebase 프로젝트의 데이터를 BigQuery로 가져왔다면, 이 코드 예시에서 BigQuery에 연결하여 데이터를 조회하고 조회 결과 데이터를 Google 클라우드 스토리지에 저장할 수 있습니다.

package cn.thinkingdata.service;

import cn.thinkingdata.bean.CloudStorageInfo;
import cn.thinkingdata.config.BQParams;
import cn.thinkingdata.util.GoogleUtil;
import com.google.api.gax.paging.Page;
import com.google.cloud.RetryOption;
import com.google.cloud.bigquery.*;
import com.google.cloud.storage.Blob;
import com.google.cloud.storage.Storage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.threeten.bp.Duration;

import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.annotation.Resource;
import java.text.SimpleDateFormat;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Calendar;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;


import static java.time.format.DateTimeFormatter.BASIC_ISO_DATE;

/**
* @Description: bigquery 쿼리로 결과를 Google 클라우드 스토리지에 저장
* @Company: ThinkingData
* @Date: 2022/1/27
* @Version: 1.0
* @Copyright: Copyright (c) 2022
*/
@Service
public class ExportToCloudStorageService {

private static final Logger logger = LoggerFactory.getLogger(ExportToCloudStorageService.class);

@Autowired
private BQParams bqParams;

@Resource(name = "StorageQueue")
private BlockingQueue<CloudStorageInfo> storageQueue;

private ThreadPoolExecutor threadPoolExecutor;
private List<String> dates;
private BigQuery bigQuery;
private Storage storage;

private String projectId;
private String dataSet;
private String region;
private String filter;
private String bucketName;
private String lower;
private String upper;
private String filePath;
private Integer export;
private SimpleDateFormat format = new SimpleDateFormat("yyyyMMdd");

@PostConstruct
public void init(){
projectId = bqParams.projectId;
dataSet = bqParams.dataset;
region = bqParams.region;
filter = bqParams.filter;
bucketName = bqParams.bucketName;
lower = bqParams.lowerbound;
upper = bqParams.upperbound;
filePath = bqParams.filePath;
export = bqParams.export;
bigQuery = GoogleUtil.getBigQuery(filePath,projectId);
storage = GoogleUtil.getStorage(filePath,projectId);
getDates();
start();
}

//기간 가져오기
private List<String> getDates(){
LocalDate lowerbound = LocalDate.parse(lower.replace("-",""), BASIC_ISO_DATE);
LocalDate upperbound = LocalDate.parse(upper.replace("-",""), BASIC_ISO_DATE);
if(upperbound.isBefore(lowerbound)){
logger.error("=====시간 상한은 시간 하한보다 이전일 수 없습니다=====");
System.exit(-1);
}
LocalDate temp = lowerbound;
LocalDate upperboundPlus = upperbound.plusDays(1);
dates = new ArrayList<>();
while (temp.isBefore(upperboundPlus)){
dates.add(temp.format(BASIC_ISO_DATE));
temp = temp.plusDays(1);
}
return dates;
}

//스레드 풀을 시작하여 데이터를 조회하고 Google 클라우드 스토리지에 저장
private void start(){
logger.info("=====데이터 조회·내보내기·저장 처리 시작======");
AtomicInteger threadNum = new AtomicInteger(0);
threadPoolExecutor = new ThreadPoolExecutor(export,export, 10,TimeUnit.SECONDS,new ArrayBlockingQueue<>(180),
r -> {
Thread thread = new Thread(r,"export-threadPool-"+threadNum.incrementAndGet());
return thread;
},new ThreadPoolExecutor.CallerRunsPolicy());
threadPoolExecutor.allowCoreThreadTimeOut(true);
//주의: 여기서는 내보내기 스레드 수를 latch 카운트로 사용할 수 없음. 7일치 데이터를 처리하는데 처리 스레드가 1개이면 포이즌 필 마커가 미리 들어가게 됨
final CountDownLatch latch = new CountDownLatch(dates.size());
for (String date : dates) {
threadPoolExecutor.submit(new ExportTask(date,latch));
}
threadPoolExecutor.submit(new EndTask(latch));
}

@PreDestroy
public void destroy(){
while (true){
if(threadPoolExecutor.getActiveCount() == 0){
logger.info("=====데이터 조회·내보내기 서비스 종료 시작=====");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("=====데이터 조회·내보내기 서비스 종료 중 예외 발생!=====",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("=====데이터 조회·내보내기 서비스 종료 성공=====");
break;
}
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}

//데이터 내보내기 스레드(bigquery에서 데이터를 조회하여 임시 테이블로 내보냄)
private class ExportTask implements Runnable {
private final String date;
private final CountDownLatch latch;
private final String jobId;

public ExportTask(String date, CountDownLatch latch) {
this.date = date;
this.latch = latch;
this.jobId = UUID.randomUUID().toString();
}

@Override
public void run() {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
try {
logger.info("=====데이터 조회 스레드: " + Thread.currentThread().getName() + " =====");
logger.info("=====당일 파티션 테이블 데이터 처리 시작=====");
dealCurrentDay();
logger.info("=====당일 파티션 테이블 데이터 처리 완료=====");
logger.info("=====1-7일 전 증분 데이터 처리 시작=====");
for (int i = 1; i <= 7; i++) {
deal1to7DayBefore(i);
}
logger.info("===== 1-7일 전 증분 데이터 처리 완료=====");
logger.info("=====데이터 조회 스레드: " + Thread.currentThread().getName() + " 처리 성공=====");
} finally {
latch.countDown();
}
}

//당일 파티션 테이블 데이터 처리
private void dealCurrentDay() {
try {
//먼저 당일 파티션 테이블 데이터용 상태 테이블 생성
String currentDayStateTable = String.format("ta_export_%s_events_%s_%s", region, date, date);
logger.info(String.format("%s 당일 상태 테이블:%s", date, currentDayStateTable));
//당일 상태 데이터를 조회하여 상태 테이블에 삽입
TableId tableId = TableId.of(projectId, dataSet, currentDayStateTable);
String sql = String.format("select * from `%s.%s.events_%s`", projectId, dataSet, date) + filter;
logger.info(String.format("%s 상태 데이터 조회 쿼리:%s", date, sql));
QueryJobConfiguration queryJobConfiguration = QueryJobConfiguration.newBuilder(sql).setUseLegacySql(false).setDestinationTable(tableId).build();
bigQuery.query(queryJobConfiguration);
Thread.sleep(1000);
//상태 테이블의 데이터를 google 클라우드에 저장
String googlePath = String.format("gs://%s/ta_export/%s-*.json.gz", bucketName, currentDayStateTable);
ExtractJobConfiguration extractJobConfiguration = ExtractJobConfiguration.newBuilder(tableId, googlePath).setCompression("gzip").setFormat(FormatOptions.json().getType()).build();
Job job = bigQuery.create(JobInfo.of(extractJobConfiguration)).waitFor(RetryOption.totalTimeout(Duration.ofMinutes(15)));
if (job.getStatus().getError() != null) {
logger.error(String.format("%s 당일 상태 테이블 데이터를 google 클라우드로 가져오기 실패: %s", date, job.getStatus().getError().getMessage()));
}
logger.info(String.format("%s 당일 상태 테이블 데이터를 google 클라우드로 가져오기 성공!", date));
Thread.sleep(1000);
//google 클라우드의 데이터를 스토리지 큐에 넣기
String prefix = String.format("ta_export/%s-", currentDayStateTable);
Page<Blob> list = storage.list(bucketName, Storage.BlobListOption.currentDirectory(), Storage.BlobListOption.prefix(prefix));
Schema schema = bigQuery.getTable(dataSet, currentDayStateTable).getDefinition().getSchema();
for (Blob blob : list.iterateAll()) {
CloudStorageInfo cloudStorageInfo = new CloudStorageInfo();
cloudStorageInfo.setBlob(blob);
cloudStorageInfo.setDate(date);
cloudStorageInfo.setEnd(false);
cloudStorageInfo.setJobId(jobId);
cloudStorageInfo.setSchema(schema);
cloudStorageInfo.setMd5(blob.getMd5ToHexString());
storageQueue.put(cloudStorageInfo);
}
} catch (Exception e) {
logger.error(String.format("당일 %s 파티션 테이블 데이터 처리 실패:%s",date,e.getMessage()));
System.exit(-1);
}
}

//1-7일 전 원본 테이블의 오늘 증분 데이터 처리(어제의 상태 테이블과 비교 필요)
private void deal1to7DayBefore(int i) {
String partitionStr = "";
String lastStr = "";
try {
Calendar partition = Calendar.getInstance();
partition.add(Calendar.DAY_OF_MONTH, -3-i);
partitionStr = format.format(partition.getTime());
Calendar last = Calendar.getInstance();
last.add(Calendar.DAY_OF_MONTH, -4);
lastStr = format.format(last.getTime());

//먼저 (-3-i)일 전 파티션 상태 테이블이 있는지 확인하고, 없으면 처리하지 않음. 프로그램 시작 시 사용됨
String lastDayStateTable = String.format("ta_export_%s_events_%s_%s", region, partitionStr, lastStr);
Table table = bigQuery.getTable(TableId.of(projectId, dataSet, lastDayStateTable));
if (table != null) {
logger.info(String.format("어제 %s 파티션 테이블의 상태 테이블이 있으므로 증분 처리 가능", partitionStr));
//원본 테이블의 오늘 상태 테이블 생성
String currentDayStateTable = String.format("ta_export_%s_events_%s_%s", region, partitionStr, date);
logger.info(String.format("%s의 %s 시점 상태 테이블:%s", partitionStr, date, currentDayStateTable));
//원본 테이블 데이터를 조회하여 현재 상태 테이블에 넣기
TableId tableId = TableId.of(projectId, dataSet, currentDayStateTable);
String sql_select = String.format("select * from `%s.%s.events_%s`", projectId, dataSet, partitionStr) + filter;
logger.info(String.format("%s 상태 데이터 조회 쿼리:%s", partitionStr, sql_select));
QueryJobConfiguration query_select = QueryJobConfiguration.newBuilder(sql_select).setUseLegacySql(false).setDestinationTable(tableId).build();
bigQuery.query(query_select);
Thread.sleep(1000);
//결과 테이블은 사용 전에 비움(이 테이블은 재사용됨)
String dest = String.format("%s.%s.%s_result_response",projectId,dataSet,region);
String sql_truncate = String.format("truncate table `%s`",dest);
logger.info(String.format("결과 테이블 비우기 쿼리:%s", sql_truncate));
QueryJobConfiguration query_truncate = QueryJobConfiguration.newBuilder(sql_truncate).setUseLegacySql(false).build();
bigQuery.query(query_truncate);
Thread.sleep(1000);
//오늘과 어제의 상태 테이블을 비교하여 증분 데이터를 걸러내고 결과 테이블에 삽입
String latest = String.format("%s.%s.ta_export_%s_events_%s_%s",projectId,dataSet,region,partitionStr,lastStr);
String current = String.format("%s.%s.ta_export_%s_events_%s_%s",projectId,dataSet,region,partitionStr,date);
String sql_insert = String.format("insert into `%s`\n" +
"select\n" +
"b_event_date,\n" +
"b_event_timestamp,\n" +
"b_event_name,\n" +
"b_event_params,\n" +
"b_event_previous_timestamp,\n" +
"b_event_value_in_usd,\n" +
"b_event_bundle_sequence_id,\n" +
"b_event_server_timestamp_offset,\n" +
"b_user_id,\n" +
"b_user_pseudo_id,\n" +
"b_privacy_info,\n" +
"b_user_properties,\n" +
"b_user_first_touch_timestamp,\n" +
"b_user_ltv,\n" +
"b_device,\n" +
"b_geo,\n" +
"b_app_info,\n" +
"b_traffic_source,\n" +
"b_stream_id,\n" +
"b_platform,\n" +
"b_event_dimensions,\n" +
"b_ecommerce,\n" +
"b_items\n" +
"from\n" +
"(select\n" +
"a.event_date a_event_date,\n" +
"a.event_timestamp a_event_timestamp,\n" +
"a.event_name a_event_name,\n" +
"a.event_params a_event_params,\n" +
"a.event_previous_timestamp a_event_previous_timestamp,\n" +
"a.event_value_in_usd a_event_value_in_usd,\n" +
"a.event_bundle_sequence_id a_event_bundle_sequence_id,\n" +
"a.event_server_timestamp_offset a_event_server_timestamp_offset,\n" +
"a.user_id a_user_id,\n" +
"a.user_pseudo_id a_user_pseudo_id,\n" +
"a.privacy_info a_privacy_info,\n" +
"a.user_properties a_user_properties,\n" +
"a.user_first_touch_timestamp a_user_first_touch_timestamp,\n" +
"a.user_ltv a_user_ltv,\n" +
"a.device a_device,\n" +
"a.geo a_geo,\n" +
"a.app_info a_app_info,\n" +
"a.traffic_source a_traffic_source,\n" +
"a.stream_id a_stream_id,\n" +
"a.platform a_platform,\n" +
"a.event_dimensions a_event_dimensions,\n" +
"a.ecommerce a_ecommerce,\n" +
"a.items a_items,\n" +
"b.event_date b_event_date,\n" +
"b.event_timestamp b_event_timestamp,\n" +
"b.event_name b_event_name,\n" +
"b.event_params b_event_params,\n" +
"b.event_previous_timestamp b_event_previous_timestamp,\n" +
"b.event_value_in_usd b_event_value_in_usd,\n" +
"b.event_bundle_sequence_id b_event_bundle_sequence_id,\n" +
"b.event_server_timestamp_offset b_event_server_timestamp_offset,\n" +
"b.user_id b_user_id,\n" +
"b.user_pseudo_id b_user_pseudo_id,\n" +
"b.privacy_info b_privacy_info,\n" +
"b.user_properties b_user_properties,\n" +
"b.user_first_touch_timestamp b_user_first_touch_timestamp,\n" +
"b.user_ltv b_user_ltv,\n" +
"b.device b_device,\n" +
"b.geo b_geo,\n" +
"b.app_info b_app_info,\n" +
"b.traffic_source b_traffic_source,\n" +
"b.stream_id b_stream_id,\n" +
"b.platform b_platform,\n" +
"b.event_dimensions b_event_dimensions,\n" +
"b.ecommerce b_ecommerce,\n" +
"b.items b_items\n" +
"from `%s` a\n" +
"right outer join `%s` b\n" +
"on a.event_name=b.event_name and a.event_timestamp=b.event_timestamp and a.device.advertising_id=b.device.advertising_id) t\n" +
"where t.b_event_name is not null and t.b_event_timestamp is not null and t.b_device.advertising_id is not null and\n" +
"t.a_event_name is null and t.a_event_timestamp is null and t.a_device.advertising_id is null", dest, latest, current);
logger.info(String.format("%s 증분 데이터 조회 쿼리:%s", partitionStr, sql_insert));
QueryJobConfiguration query_insert = QueryJobConfiguration.newBuilder(sql_insert).setUseLegacySql(false).build();
bigQuery.query(query_insert);
Thread.sleep(1000);
//결과 테이블의 데이터(증분)를 google 클라우드에 저장
String destTable = String.format("%s_result_response",region);
TableId destTableId = TableId.of(projectId, dataSet, destTable);
String googlePath = String.format("gs://%s/ta_export/%s-*.json.gz", bucketName, destTable);
ExtractJobConfiguration extractJobConfiguration = ExtractJobConfiguration.newBuilder(destTableId, googlePath).setCompression("gzip").setFormat(FormatOptions.json().getType()).build();
Job job = bigQuery.create(JobInfo.of(extractJobConfiguration)).waitFor(RetryOption.totalTimeout(Duration.ofMinutes(15)));
if (job.getStatus().getError() != null) {
logger.error(String.format("%s 증분 데이터를 google 클라우드로 가져오기 실패: %s", partitionStr, job.getStatus().getError().getMessage()));
}
logger.info(String.format("%s 증분 데이터를 google 클라우드로 가져오기 성공!", partitionStr));
Thread.sleep(1000);
//증분 데이터를 걸러낸 후에는 어제의 상태 테이블을 삭제해야 함. 상태 테이블은 매일 업데이트되기 때문임. 매번 비교하는 대상은 특정 날짜 파티션 테이블의 오늘 상태 테이블과 어제까지의 최신 상태 테이블임
//오늘 비교를 마치면 어제의 상태 테이블은 삭제해도 됨. 내일은 내일의 상태 테이블과 오늘의 상태 테이블을 비교하므로 어제의 상태 테이블은 더 이상 필요 없음.
boolean delete = bigQuery.delete(TableId.of(projectId, dataSet, lastDayStateTable));
if (delete) {
logger.info(String.format("%s의 %s 시점 상태 테이블 삭제 성공!", partitionStr, lastStr));
} else {
logger.error(String.format("%s의 %s 시점 상태 테이블 삭제 실패!!!", partitionStr, lastStr));
}
Thread.sleep(1000);
//google 클라우드의 데이터를 스토리지 큐에 넣기
String prefix = String.format("ta_export/%s-", destTable);
Page<Blob> list = storage.list(bucketName, Storage.BlobListOption.currentDirectory(), Storage.BlobListOption.prefix(prefix));
Schema schema = bigQuery.getTable(TableId.of(projectId,dataSet,destTable)).getDefinition().getSchema();
for (Blob blob : list.iterateAll()) {
CloudStorageInfo cloudStorageInfo = new CloudStorageInfo();
cloudStorageInfo.setBlob(blob);
cloudStorageInfo.setDate(partitionStr);
cloudStorageInfo.setEnd(false);
cloudStorageInfo.setJobId(jobId);
cloudStorageInfo.setSchema(schema);
cloudStorageInfo.setMd5(blob.getMd5ToHexString());
storageQueue.put(cloudStorageInfo);
}
} else {
logger.info(String.format("어제 %s 파티션 테이블의 상태 테이블이 없으므로 증분 처리 없음", partitionStr));
}
} catch (Exception e) {
logger.error(String.format("%s 날짜의 파티션 데이터를 오늘(%s) 처리하는 데 실패:%s",partitionStr,date,e.getMessage()));
System.exit(-1);
}
}
}

//포이즌 필 데이터 마커를 스토리지 큐에 넣기
private class EndTask implements Runnable{

private CountDownLatch latch;

public EndTask(CountDownLatch latch){
this.latch = latch;
}

@Override
public void run() {
try {
latch.await();
} catch (InterruptedException e) {
e.printStackTrace();
}
logger.info("=====스토리지 큐에 포이즌 필 마커 넣기 시작=====");
CloudStorageInfo cloudStorageInfo = new CloudStorageInfo();
cloudStorageInfo.setEnd(true);
try {
storageQueue.put(cloudStorageInfo);
} catch (InterruptedException e) {
logger.error("=====스토리지 큐 포이즌 필 마커 넣기 실패=====",e);
}
logger.info("=====스토리지 큐 포이즌 필 마커 넣기 성공=====");
}
}
}

이 코드에는 BigQuery 연결 및 인증, BigQuery 테이블 가져오기, sql 쿼리를 실행하여 데이터를 조회하고 결과를 임시 테이블에 넣기, 임시 테이블의 데이터를 클라우드 스토리지에 덤프하여 파일로 만들기 등의 내용이 포함되어 있습니다.

2.2.2 클라우드 스토리지에서 로컬 머신으로 파일 다운로드​

package cn.thinkingdata.service;

import cn.thinkingdata.bean.CloudStorageInfo;
import cn.thinkingdata.bean.LocalStorageInfo;
import cn.thinkingdata.config.BQParams;
import cn.thinkingdata.util.RetryerUtil;
import com.github.rholder.retry.RetryException;
import com.github.rholder.retry.Retryer;
import com.google.cloud.storage.Blob;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.annotation.Resource;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/**
* @Description: Google 클라우드 스토리지의 데이터를 로컬 스토리지로 가져오기
* @Company: ThinkingData
* @Date: 2022/1/29
* @Version: 1.0
* @Copyright: Copyright (c) 2022
*/
@Service
public class FetchDataToLocalService {

private final static Logger logger = LoggerFactory.getLogger(FetchDataToLocalService.class);

@Autowired
private BQParams bqParams;

@Resource(name = "StorageQueue")
private BlockingQueue<CloudStorageInfo> storageQueue;

@Resource(name = "FetchQueue")
private BlockingQueue<LocalStorageInfo> fetchQueue;

private ThreadPoolExecutor threadPoolExecutor;
private String localDir;
private Integer fetch;

private final Retryer retryer = RetryerUtil.initRetryerByIncTimes(5, 10L, 30L);

@PostConstruct
public void init(){
localDir = bqParams.localDir;
fetch = bqParams.fetch;

//로컬 저장 경로 생성
if(!Files.exists(Paths.get(localDir))){
try {
Files.createDirectories(Paths.get(localDir));
} catch (IOException e) {
logger.error("==========로컬 저장 경로 생성 실패==========",e);
return;
}
}
start();
}

//스토리지 큐에서 데이터를 꺼내 로컬 데이터 형식으로 변환한 후 다운로드 큐에 넣기.
private void start(){
logger.info("==========스토리지 큐에서 데이터 가져오기 시작==========");
AtomicInteger threadNum = new AtomicInteger(0);
threadPoolExecutor = new ThreadPoolExecutor(fetch, fetch, 10, TimeUnit.SECONDS, new ArrayBlockingQueue<>(180),
r -> {
Thread thread = new Thread(r, "fetch-threadPool-" + threadNum.incrementAndGet());
return thread;
}, new ThreadPoolExecutor.CallerRunsPolicy());
threadPoolExecutor.allowCoreThreadTimeOut(true);

final CountDownLatch latch = new CountDownLatch(fetch);
for (int i = 0; i < fetch; i++) {
threadPoolExecutor.submit(new FetchTask(latch));
}
threadPoolExecutor.submit(new EndTask(latch));
}

@PreDestroy
public void destroy(){
while (true){
if(threadPoolExecutor.getActiveCount() == 0){
logger.info("==========데이터 다운로드 서비스 종료 시작==========");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("==========데이터 다운로드 서비스 종료 중 예외 발생!==========",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("==========데이터 다운로드 서비스 종료 성공==========");
break;
}
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}

private class FetchTask implements Runnable{

private final CountDownLatch latch;

public FetchTask(CountDownLatch latch){
this.latch = latch;
}

@Override
public void run() {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
logger.info("==========데이터 다운로드 스레드: " + Thread.currentThread().getName()+"==========");
while (true){
try {
//스토리지 큐에서 데이터 가져오기
CloudStorageInfo cloudStorageInfo = storageQueue.take();
if(cloudStorageInfo == null){
Thread.sleep(1000);
continue;
}
//포이즌 필 데이터가 아닌 경우의 처리
if(!cloudStorageInfo.isEnd()){
Blob blob = cloudStorageInfo.getBlob();
String name = blob.getName().replace("/", "");
Path localDirPath = Paths.get(localDir).resolve(Paths.get(cloudStorageInfo.getDate()));
if(!Files.exists(localDirPath)){
Files.createDirectories(localDirPath);
}
Path localPath = localDirPath.resolve(Paths.get(name));

retryer.call(() -> {
try {
blob.downloadTo(localPath);
logger.info("==========" + blob.getName()+ " Google 클라우드 스토리지에서 로컬로 다운로드 성공==========");
} catch (Exception e) {
// 여기서 Connection closed prematurely: bytesRead = ?, Content-Length = ? 예외가 발생할 수 있음
// 원인은 Google 네트워크 변동과 관련이 있을 수 있음. 여기서 다운로드에 실패하면 파일이 불완전해지고, 이후 데이터 처리 스레드가 이 불완전한 파일을 처리할 때
// 데이터 처리 스레드가 실패함. 연쇄 반응이 생기므로 향후 최적화 방법을 검토해야 하며, 기술 지원 담당자가 고려해 볼 수 있음
logger.error("==========" + blob.getName() +" 다운로드 실패==========",e);
}
return null;
}
);

//로컬 저장 형식 데이터 생성
LocalStorageInfo localStorageInfo = new LocalStorageInfo();
localStorageInfo.setPath(localPath);
localStorageInfo.setEnd(false);
localStorageInfo.setJobId(cloudStorageInfo.getJobId());
localStorageInfo.setSchema(cloudStorageInfo.getSchema());
localStorageInfo.setMd5(cloudStorageInfo.getMd5());

fetchQueue.put(localStorageInfo);
boolean delete = blob.delete();
if(delete){
logger.info("==========" + blob.getName()+ " Google 클라우드 스토리지에서 삭제 성공==========");
}else{
logger.error("==========" + blob.getName() +" Google 클라우드 스토리지에서 삭제 실패==========");
}
}else{
//포이즌 필 데이터 처리(마지막에 처리됨)
//포이즌 필 데이터를 만나면 현재 다운로드 스레드를 중지하고 포이즌 필 데이터를 스토리지 큐에 다시 넣음
//멀티스레드인 경우 다른 스레드도 스토리지 큐에서 포이즌 필을 가져와 중지함. 최종적으로 모든 스레드가 중지됨
storageQueue.put(cloudStorageInfo);
latch.countDown();
break;
}
} catch (InterruptedException | IOException | ExecutionException | RetryException e) {
logger.error("==========데이터 다운로드 스레드:" + Thread.currentThread().getName()+" 처리 실패==========",e);
}
}//end while
logger.info("==========데이터 다운로드 스레드:" + Thread.currentThread().getName()+" 처리 성공==========");
}
}

private class EndTask implements Runnable{

private final CountDownLatch latch;

public EndTask(CountDownLatch latch){
this.latch = latch;
}

@Override
public void run() {
try {
latch.await();
} catch (InterruptedException e) {
e.printStackTrace();
}
logger.info("==========다운로드 큐에 포이즌 필 마커 넣기 시작==========");
LocalStorageInfo localStorageInfo = new LocalStorageInfo();
localStorageInfo.setEnd(true);
try {
fetchQueue.put(localStorageInfo);
} catch (InterruptedException e) {
logger.error("==========다운로드 큐 포이즌 필 마커 넣기 실패==========",e);
}
logger.info("==========다운로드 큐 포이즌 필 마커 넣기 성공==========");
}
}
}

이 코드는 Google 클라우드 스토리지에서 데이터 파일을 로컬 머신으로 다운로드하고, 다운로드가 끝나면 클라우드 스토리지의 임시 데이터 파일을 삭제하는 방법 등을 보여 줍니다.

2.2.3 로컬 파일 내용을 읽어 처리한 후 AE로 전송​

package cn.thinkingdata.service;

import cn.thinkingdata.bean.LocalStorageInfo;
import cn.thinkingdata.config.BQParams;
import cn.thinkingdata.util.CompressUtil;
import cn.thinkingdata.util.DataParseUtil;
import cn.thinkingdata.util.HttpRequestUtil;
import cn.thinkingdata.util.RetryerUtil;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.github.rholder.retry.RetryException;
import com.github.rholder.retry.Retryer;
import com.google.cloud.bigquery.FieldList;
import com.google.cloud.bigquery.Schema;
import com.google.common.util.concurrent.RateLimiter;
import org.apache.http.HttpStatus;
import org.apache.http.client.methods.CloseableHttpResponse;
import org.apache.http.client.methods.HttpPost;
import org.apache.http.entity.ByteArrayEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.util.EntityUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.annotation.Resource;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.zip.GZIPInputStream;

/**
* @Description: 다운로드 큐에서 데이터를 가져와 처리한 후 AE로 전송
* @Company: ThinkingData
* @Date: 2022/1/30
* @Version: 1.0
* @Copyright: Copyright (c) 2022
*/
@Service
public class ProcessAndSendDataService {

private static final Logger logger = LoggerFactory.getLogger(ProcessAndSendDataService.class);

@Autowired
private BQParams bqParams;

private ThreadPoolExecutor threadPoolExecutor;

@Resource(name = "FetchQueue")
private BlockingQueue<LocalStorageInfo> fetchQueue;

private String receiver;
private String appid;
private Integer process;
private Integer batchSize;

private final Retryer retryer = RetryerUtil.initNeverStopLockRetryer();

private static final RateLimiter rateLimiter = RateLimiter.create(2000);

@PostConstruct
public void init(){
receiver = bqParams.receiver;
appid = bqParams.appid;
process = bqParams.process;
batchSize = bqParams.batchSize;
start();
}

//다운로드 큐에서 데이터를 가져와 처리한 후 AE로 전송
private void start(){
logger.info("==========다운로드 큐에서 데이터 가져오기 시작==========");
AtomicInteger threadNum = new AtomicInteger(0);
threadPoolExecutor = new ThreadPoolExecutor(process, process, 10, TimeUnit.SECONDS, new ArrayBlockingQueue<>(180),
r -> {
Thread thread = new Thread(r, "process-threadPool-" + threadNum.incrementAndGet());
return thread;
}, new ThreadPoolExecutor.CallerRunsPolicy());
threadPoolExecutor.allowCoreThreadTimeOut(true);

for (int i = 0; i < process; i++) {
threadPoolExecutor.submit(new SendTask());
}
}

@PreDestroy
public void destroy(){
while (true){
if(threadPoolExecutor.getActiveCount() == 0){
logger.info("==========데이터 처리·전송 서비스 종료 시작==========");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("==========데이터 처리·전송 서비스 종료 중 예외 발생==========",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("==========데이터 처리·전송 서비스 종료 성공==========");
break;
}
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}

private class SendTask implements Runnable{

@Override
public void run() {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
logger.info("==========데이터 처리 스레드: "+Thread.currentThread().getName()+"==========");
while (true){
try {
LocalStorageInfo localStorageInfo = fetchQueue.take();
if(localStorageInfo == null){
Thread.sleep(1000);
continue;
}
JSONArray jsonArray = new JSONArray();
//포이즌 필 데이터가 아닌 경우의 처리
if(!localStorageInfo.isEnd()){
Schema schema = localStorageInfo.getSchema();
FieldList fieldList = schema.getFields();
Path path = localStorageInfo.getPath();
GZIPInputStream gzipInputStream = new GZIPInputStream(Files.newInputStream(path));
InputStreamReader inputStreamReader = new InputStreamReader(gzipInputStream);
BufferedReader bufferedReader = new BufferedReader(inputStreamReader);
String line;
int count = 0;
while ((line = bufferedReader.readLine()) != null ){
Map<String,String> fieldInfo = new HashMap<>();
JSONObject jsonObject = JSONObject.parseObject(line);
DataParseUtil.getAllFields(fieldList, jsonObject, fieldInfo);
JSONObject result = new JSONObject();
DataParseUtil.expand(jsonObject,fieldInfo,result,"");
JSONObject taObject = DataParseUtil.getTAObject(result, localStorageInfo.getJobId());
//속도 제한
while (true){
//토큰 1개 획득
boolean acquire = rateLimiter.tryAcquire();
//토큰을 획득한 경우에만 데이터를 추가하고, 그렇지 않으면 토큰을 얻을 때까지 블로킹하여 대기함
if(acquire){
jsonArray.add(taObject);
count++;
//배치 크기에 도달해야 데이터를 전송함
if(count % batchSize == 0){
sendToReceiver(jsonArray);
jsonArray.clear();
}
break;
}
}
}
if(!jsonArray.isEmpty()){
sendToReceiver(jsonArray);
}
gzipInputStream.close();
inputStreamReader.close();
bufferedReader.close();
logger.info("==========로컬 파일:"+path+"에서 총 "+count+"건의 데이터를 처리함==========");
Files.deleteIfExists(path);
logger.info("==========로컬 파일 삭제 성공:"+path+"==========");
}else{
//포이즌 필 데이터 처리
//포이즌 필 데이터가 있으면 다운로드 큐에 다시 넣고 현재 스레드를 중지함
//다른 스레드도 큐에서 이 데이터를 가져오면 중지하며, 최종적으로 모두 중지됨
fetchQueue.put(localStorageInfo);
break;
}
} catch (InterruptedException | IOException | ExecutionException | RetryException e) {
//여기서 java.io.EOFException: Unexpected end of ZLIB 예외가 발생할 수 있으며, 스레드가 종료되어 파일을 처리하지 못하고 삭제되지도 않음. 향후 어떻게 최적화할지 검토 필요?
//예: 새 스레드를 띄워 다시 소비? 기술 지원 담당자가 고려해 볼 수 있음
logger.error("==========데이터 처리 스레드:"+Thread.currentThread().getName()+" 예외 발생==========",e);
}
} //end while
logger.info("==========데이터 처리 스레드:"+Thread.currentThread().getName()+" 정상 종료==========");
}

private void sendToReceiver(JSONArray jsonArray) throws ExecutionException, RetryException {
retryer.call(() -> {
byte[] dataBytes = CompressUtil.gzipCompress(jsonArray.toJSONString().getBytes(StandardCharsets.UTF_8));
if (dataBytes != null) {
CloseableHttpClient httpClient = HttpRequestUtil.getHttpClient();
HttpPost post = new HttpPost(receiver);
post.setHeader("compress","gzip");
post.addHeader("appid", appid);
post.addHeader("user-agent","datax-1.0");
post.setEntity(new ByteArrayEntity(dataBytes));
CloseableHttpResponse httpResponse = httpClient.execute(post);
int statusCode = httpResponse.getStatusLine().getStatusCode();
EntityUtils.consumeQuietly(httpResponse.getEntity());
if (statusCode != HttpStatus.SC_OK) {
throw new Exception("error response status code: " + statusCode);
}
}
return null;
});
}
}
}

이 코드는 로컬 파일 읽기, 데이터 파싱·처리, AE 플랫폼으로 전송 등의 내용을 보여 줍니다.

3. 프로젝트 배포 및 시작​

이 코드 프로젝트는 springboot로 개발되었습니다. 프로젝트 패키징은 springboot 방식을 따르며, 스크립트 파일로 시작합니다. 스크립트 파일의 설정 예시는 다음과 같습니다

#!/usr/bin/env bash

source ~/.bash_profile
export LANG="en_US.UTF-8"

start_day=`date -d "-3 days" "+%Y-%m-%d"`
end_day=`date -d "-3 days" "+%Y-%m-%d"`

cd `dirname $0`

JAVA_JVM_ARGS="-Xmx4096m -Xms512m -XX:+UseG1GC -XX:+DisableExplicitGC -XX:+UseGCOverheadLimit"

jar="ta-data-main-1.0.0.jar"

nohup java -server ${JAVA_JVM_ARGS} \
-jar -Dloader.path=../lib -Dlogname=eur-your-app ../jars/$jar \
--bigquery.filePath=../json/your-service-account.json \
--bigquery.date.lowerbound=$start_day \
--bigquery.date.upperbound=$end_day \
--spring.config.location=../conf/application.yml \
--spring.profiles.active=eur-your-app \
> /dev/null 2>&1 < /dev/null &

시작 스크립트에서 --로 시작하는 설정 파라미터는 무시해도 됩니다. 이 파라미터는 고객, 프로젝트, 요구 사항에 따라 다르게 설정하며 참고용입니다.

이 문서가 도움이 되었나요?