メインコンテンツまでスキップ

Firebase-BigQuery-GCSデータ移行の技術プラン

最終更新 2026/10/07

最終更新日:2022-04-13

概要​

Firebaseは統合ツールで、さまざまなアプリケーションと簡単に統合できます。公式ドキュメントは次のとおりです

https://firebase.google.com/docs

プラットフォーム別にFirebaseのドキュメントを見る

プロダクト別に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. サンプルプロジェクト​

サンプルプロジェクトには現在公開のダウンロードURLがありません。入手については、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: Stores the results of BigQuery queries in Google Cloud Storage
* @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();
}

//Get the time range
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("=====The upper time bound cannot be earlier than the lower time bound=====");
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;
}

//Start the thread pool to query data and store it in Google Cloud Storage
private void start(){
logger.info("=====Start querying, exporting, and storing data======");
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);
//Note: the number of export threads cannot be used as the latch count here. If 7 days of data are processed with only 1 processing thread, the poison pill marker would be added too early
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("=====Start shutting down the data query and export service=====");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("=====Error while shutting down the data query and export service!=====",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("=====Data query and export service shut down successfully=====");
break;
}
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}

//Data export thread (queries data from BigQuery and exports it to a temporary table)
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("=====Data query thread: " + Thread.currentThread().getName() + " =====");
logger.info("=====Start processing the data in the current day's partition table=====");
dealCurrentDay();
logger.info("=====Finished processing the current day's partition table data=====");
logger.info("=====Start processing incremental data from 1-7 days ago=====");
for (int i = 1; i <= 7; i++) {
deal1to7DayBefore(i);
}
logger.info("===== Finished processing incremental data from 1-7 days ago=====");
logger.info("=====Data query thread: " + Thread.currentThread().getName() + " processed successfully=====");
} finally {
latch.countDown();
}
}

//Process the current day's partition table data
private void dealCurrentDay() {
try {
//First, create a state table for the current day's partition table data
String currentDayStateTable = String.format("ta_export_%s_events_%s_%s", region, date, date);
logger.info(String.format("%s current-day state table: %s", date, currentDayStateTable));
//Query the current day's state data and insert it into the state table
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("Statement to query %s state data: %s", date, sql));
QueryJobConfiguration queryJobConfiguration = QueryJobConfiguration.newBuilder(sql).setUseLegacySql(false).setDestinationTable(tableId).build();
bigQuery.query(queryJobConfiguration);
Thread.sleep(1000);
//Store the data in the state table in Google Cloud
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 failed to import current-day state table data to Google Cloud: %s", date, job.getStatus().getError().getMessage()));
}
logger.info(String.format("%s imported current-day state table data to Google Cloud successfully!", date));
Thread.sleep(1000);
//Put the data in Google Cloud into the storage queue
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("Failed to process the current day's %s partition table data: %s",date,e.getMessage()));
System.exit(-1);
}
}

//Process today's incremental data of the source tables from 1-7 days ago (compared with yesterday's state table)
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());

//First check whether the partition state table from (-3-i) days ago exists. If it does not exist, skip processing. This is used at the very start of the program
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("A state table for the %s partition table existed yesterday, so incremental processing can proceed", partitionStr));
//Generate today's state table for the source table
String currentDayStateTable = String.format("ta_export_%s_events_%s_%s", region, partitionStr, date);
logger.info(String.format("State table of %s on %s: %s", partitionStr, date, currentDayStateTable));
//Query the source table data and put it into the current state table
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("Statement to query %s state data: %s", partitionStr, sql_select));
QueryJobConfiguration query_select = QueryJobConfiguration.newBuilder(sql_select).setUseLegacySql(false).setDestinationTable(tableId).build();
bigQuery.query(query_select);
Thread.sleep(1000);
//Clear the result table before use. This table is reused
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("Statement to clear the result table: %s", sql_truncate));
QueryJobConfiguration query_truncate = QueryJobConfiguration.newBuilder(sql_truncate).setUseLegacySql(false).build();
bigQuery.query(query_truncate);
Thread.sleep(1000);
//Compare the state tables of today and yesterday, filter out the incremental data, and insert it into the result table
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("Statement to get %s incremental data: %s", partitionStr, sql_insert));
QueryJobConfiguration query_insert = QueryJobConfiguration.newBuilder(sql_insert).setUseLegacySql(false).build();
bigQuery.query(query_insert);
Thread.sleep(1000);
//Store the (incremental) data in the result table in Google Cloud
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 failed to import incremental data to Google Cloud: %s", partitionStr, job.getStatus().getError().getMessage()));
}
logger.info(String.format("%s imported incremental data to Google Cloud successfully!", partitionStr));
Thread.sleep(1000);
//After the incremental data is filtered out, delete the state table of yesterday, because state tables are updated every day. Each comparison is between the state table of a partition table for today and its latest state table as of yesterday
//After the comparison today, the state table of yesterday can be deleted. Tomorrow, the state table of tomorrow is compared with the state table of today, so the state table of yesterday is no longer needed.
boolean delete = bigQuery.delete(TableId.of(projectId, dataSet, lastDayStateTable));
if (delete) {
logger.info(String.format("State table of %s on %s deleted successfully!", partitionStr, lastStr));
} else {
logger.error(String.format("Failed to delete the state table of %s on %s!!!", partitionStr, lastStr));
}
Thread.sleep(1000);
//Put the data in Google Cloud into the storage queue
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("No state table for the %s partition table existed yesterday, so there is no incremental processing", partitionStr));
}
} catch (Exception e) {
logger.error(String.format("Failed to process the partition data of %s on %s (today): %s",partitionStr,date,e.getMessage()));
System.exit(-1);
}
}
}

//Put the poison pill marker into the storage queue
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("=====Start putting the poison pill marker into the storage queue=====");
CloudStorageInfo cloudStorageInfo = new CloudStorageInfo();
cloudStorageInfo.setEnd(true);
try {
storageQueue.put(cloudStorageInfo);
} catch (InterruptedException e) {
logger.error("=====Failed to put the poison pill marker into the storage queue=====",e);
}
logger.info("=====Poison pill marker put into the storage queue successfully=====");
}
}
}

このコードには、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: Pulls data from Google Cloud Storage to local storage
* @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;

//Create the local storage path
if(!Files.exists(Paths.get(localDir))){
try {
Files.createDirectories(Paths.get(localDir));
} catch (IOException e) {
logger.error("==========Failed to create the local storage path==========",e);
return;
}
}
start();
}

//Take data from the storage queue, convert it to the local data format, and then put it into the fetch queue.
private void start(){
logger.info("==========Start taking data from the storage queue==========");
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("==========Start shutting down the data fetch service==========");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("==========Error while shutting down the data fetch service!==========",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("==========Data fetch service shut down successfully==========");
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("==========Data fetch thread: " + Thread.currentThread().getName()+"==========");
while (true){
try {
//Take data from the storage queue
CloudStorageInfo cloudStorageInfo = storageQueue.take();
if(cloudStorageInfo == null){
Thread.sleep(1000);
continue;
}
//Handle data that is not a poison pill
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()+ " downloaded from Google Cloud Storage to local successfully==========");
} catch (Exception e) {
// A Connection closed prematurely: bytesRead = ?, Content-Length = ? exception may occur here
// The cause may be related to Google network fluctuations. If the download fails here, the file is incomplete, and when a later data processing thread processes this incomplete file,
// the data processing thread fails. This causes a chain reaction. Consider how to optimize this later. Technical support colleagues can look into it
logger.error("==========" + blob.getName() +" download failed==========",e);
}
return null;
}
);

//Create data in the local storage format
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()+ " deleted from Google Cloud Storage successfully==========");
}else{
logger.error("==========" + blob.getName() +" could not be deleted from Google Cloud Storage==========");
}
}else{
//Handle poison pill data (processed last)
//If poison pill data is encountered, stop the current fetch thread and put the poison pill data back into the storage queue
//With multiple threads, the other threads also get the poison pill from the storage queue and stop. Eventually all threads stop
storageQueue.put(cloudStorageInfo);
latch.countDown();
break;
}
} catch (InterruptedException | IOException | ExecutionException | RetryException e) {
logger.error("==========Data fetch thread: " + Thread.currentThread().getName()+" failed==========",e);
}
}//end while
logger.info("==========Data fetch thread: " + Thread.currentThread().getName()+" processed successfully==========");
}
}

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("==========Start putting the poison pill marker into the fetch queue==========");
LocalStorageInfo localStorageInfo = new LocalStorageInfo();
localStorageInfo.setEnd(true);
try {
fetchQueue.put(localStorageInfo);
} catch (InterruptedException e) {
logger.error("==========Failed to put the poison pill marker into the fetch queue==========",e);
}
logger.info("==========Poison pill marker put into the fetch queue successfully==========");
}
}
}

このコードでは、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 &

起動スクリプト内の -- で始まる設定パラメータは無視してかまいません。これはお客様、プロジェクト、要件ごとに設定するもので、参考用です。

このページは役に立ちましたか?