Firebase-BigQuery-GCS data migration technical solution
Last updated: 2022-04-13
Overview
Firebase is an integration tool that can easily be integrated with many apps. The official documentation is as follows
https://firebase.google.com/docs
Firebase documentation by platform
Firebase documentation by product
Data in a Firebase project can be exported to many destinations. Because Firebase is a Google product, many users export their data to a cloud data warehouse, that is, BigQuery, and use the powerful SQL features of BigQuery to analyze the data. So in general, Firebase + BigQuery is the choice of most developers.
Exporting data from a Firebase project to BigQuery must be done by the Firebase user. See this article in the official Firebase documentation.
https://firebase.google.com/docs/projects/bigquery-export
This document applies only to the following process: export data from a Firebase project to BigQuery, query and filter the data with the SQL query feature of BigQuery, transfer the query results to cloud storage (that is, Google Cloud Storage), download the data files from cloud storage to a local machine, and then parse and process the file data and send it to the AE platform.
The image below shows the overall data flow
1. Configuration the customer needs to provide
- The JSON file of a Service Account that has read and write permissions for Google BigQuery and Google Cloud Storage. Here's a sample configuration
{
"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>"
}
- The bucketName of Google Cloud Storage (Note that Google Cloud Storage and Google BigQuery must be in the same region). Here's a sample configuration
bucketName: "your-bucket-name"
- The Google BigQuery project name, which is the name under which a specific project is stored in BigQuery. Here's a sample configuration
projectId: "your-project-id"
- The dataset under a Google BigQuery project. A dataset is the specific storage location under a project, and it usually contains the tables that store the data. Here's a sample configuration
dataset: "analytics_123456789"
- The wildcard for data tables. If all the tables in the dataset that store a project's data have the same name prefix, you can also configure this prefix (the prefix is optional).
2. Sample project
The sample project has no public download URL yet. Contact your ThinkingAI customer success manager to get it.
2.1 Sample project structure
2.2 Code examples
The code is provided only as an example and for reference. Don't copy it as is!
This code example only shows how to write a program with the BigQuery SDK. Modify the business logic to fit your actual business!
2.2.1 Query data from BigQuery and save the results to cloud storage
If you've already imported the data of your Firebase project into BigQuery, this code example connects to BigQuery to query the data and stores the query results in Google Cloud Storage.
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=====");
}
}
}
This code covers connecting to BigQuery and authenticating, getting BigQuery tables, running SQL queries and putting the results into a temporary table, and dumping the data in the temporary table to cloud storage as files.
2.2.2 Download files from cloud storage to a local machine
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==========");
}
}
}
This code shows how to download data files from Google Cloud Storage to a local machine and delete the temporary data files from cloud storage after the download finishes.
2.2.3 Read local files, process the content, and send it to 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: Gets data from the fetch queue, processes it, and sends it to 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();
}
//Get data from the fetch queue, process it, and send it to AE
private void start(){
logger.info("==========Start getting data from the fetch queue==========");
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("==========Start shutting down the data processing and sending service==========");
threadPoolExecutor.shutdown();
if(!threadPoolExecutor.isTerminated()){
try {
threadPoolExecutor.awaitTermination(60,TimeUnit.SECONDS);
} catch (InterruptedException e) {
logger.error("==========Error while shutting down the data processing and sending service==========",e);
}
threadPoolExecutor.shutdownNow();
}
logger.info("==========Data processing and sending service shut down successfully==========");
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("==========Data processing thread: "+Thread.currentThread().getName()+"==========");
while (true){
try {
LocalStorageInfo localStorageInfo = fetchQueue.take();
if(localStorageInfo == null){
Thread.sleep(1000);
continue;
}
JSONArray jsonArray = new JSONArray();
//Handle data that is not a poison pill
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());
//Rate limiting
while (true){
//Acquire a token
boolean acquire = rateLimiter.tryAcquire();
//Add data only after a token is acquired. Otherwise, block and wait until a token is acquired
if(acquire){
jsonArray.add(taObject);
count++;
//Send data only when the batch size is reached
if(count % batchSize == 0){
sendToReceiver(jsonArray);
jsonArray.clear();
}
break;
}
}
}
if(!jsonArray.isEmpty()){
sendToReceiver(jsonArray);
}
gzipInputStream.close();
inputStreamReader.close();
bufferedReader.close();
logger.info("==========Local file: "+path+" processed a total of "+count+" records==========");
Files.deleteIfExists(path);
logger.info("==========Deleted local file successfully: "+path+"==========");
}else{
//Handle poison pill data
//If poison pill data exists, put it back into the fetch queue and stop the current thread
//Other threads that get this data from the queue also stop, until all threads have stopped
fetchQueue.put(localStorageInfo);
break;
}
} catch (InterruptedException | IOException | ExecutionException | RetryException e) {
//A java.io.EOFException: Unexpected end of ZLIB exception may occur here. The thread then ends, so the file cannot be processed and is not deleted. How should this be optimized later?
//For example, start a new thread to consume it again? Technical support colleagues can look into this
logger.error("==========Data processing thread: "+Thread.currentThread().getName()+" encountered an exception==========",e);
}
} //end while
logger.info("==========Data processing thread: "+Thread.currentThread().getName()+" finished successfully==========");
}
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;
});
}
}
}
This code shows how to read local files, parse and process the data, and send it to the AE platform.
3. Project deployment and startup
This project is developed with Spring Boot. Package it in the Spring Boot style and start it with a script file. Here's a sample configuration of the script file
#!/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 &
The configuration parameters that start with -- in the startup script can be ignored. Set them according to the customer, project, and requirements. For reference only.

