Firebase-BigQuery-GCS 数据迁移技术方案
最近更新时间:2022-04-13
简介
Firebase 是一个集成工具,可以方便的与一些应用进行集成。官方文档如下
https://firebase.google.com/docs
按平台查看 Firebase 文档
按产品查看 Firebase 文档
Firebase 项目中的数据可以导出到很多地方。由于是谷歌旗下的,而且很多的使用者会将数据导出到云存储仓库,也就是 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. 客户需要提供的配置
- 具有 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>"
}
- Google Cloud Storage 的 bucketName(注意 Google Cloud Storage 和 Google BigQuery 需要在一个 region),下面是一份配置样例
bucketName: "your-bucket-name"
- Google BigQuery 的项目名称,代表某一个具体的项目存储在 BigQuery 上的名字。下面是一份配置样例
projectId: "your-project-id"
- Google BigQuery 某一个项目下的数据集,数据集代表某一个项目下的具体存放位置。数据集下面通常就是存储数据的表了。下面是一份配置样例
dataset: "analytics_123456789"
- 数据表的通配符。如果一个项目下的所有数据在数据集中存储的所有表名称都有相同的前缀,那么也可以配置一下这个前缀(该前缀非必须)。
2. 样例工程
样例工程暂无公共下载地址,请联系 ThinkingAI 客户成功经理获取。
2.1 样例工程结构
2.2 代码示例
代码仅提供示例和思路,切勿直接照搬!
此代码案例仅展示如何使用 BigQuery 的 sdk 进行程序的编写。具体业务逻辑需要根据自己的实际业务进行修改!
2.2.1 从 BigQuery 中查询数据并将结果存入云存储
在此之前,如果有将 Firebase 项目中的数据导入到 BigQuery 中,那么在此代码示例中就可以连接 BigQuery 进行数据的查询,并且将查询的结果数据存储到谷歌云存储上。
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 = ? 异常
// 原因可能跟谷歌的网络波动有关,如果这里下载失败,会导致文件不完整,那么后面的数据处理线程处理到该不完整文件时
// 会导致数据处理线程失败.会有连锁反应,看后续如何优化,技术支持的同事可以考虑一下
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("==========抓取队列毒药标记放入成功==========");
}
}
}
此代码展示了如何从谷歌云存储上将数据文件下载到本地机器,下载结束后删除云存储上的临时数据文件等内容。
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){
//获取一个令牌
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 &
启动脚本中的 -- 开头的配置参数可以忽略,这个根据不同的客户、不同的项目、不同的需求来进行设定,仅供参考。

