简易数据仓库搭建
·
日志服务器
Flume
Kafka
Flume->HDFS
Hive离线数仓搭建
ODS原始数据层
1:创建外部表
2:
startlog



eventlog

编写加载数据脚本-追求每隔一段时间将数据导入到hive

DWD明细数据层
启动日志展开


创建启动日志和事件日志的base表
##创建启动日志基础表
drop table if exists dwd_start_log;
create external table dwd_base_start_log(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`event_name` string,
`event_json` string,
`server_time` string
)
partitioned by(dt string)
stored as parquet
location '/warehouse/voicebar/dwd/dwd_base_start_log';
## 创建事件基础表
drop table if exists dwd_event_log;
create external table dwd_base_event_log(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`event_name` string,
`event_json` string,
`server_time` string
)
partitioned by(dt string)
stored as parquet
location '/warehouse/voicebar/dwd/dwd_base_event_log';


自定义UDF函数解析公共字段
public class BaseFieldUDF extends UDF {
public static void main(String[] args) {
String line = "1585989597235|{\"cm\":{\"ln\":\"-87.6\",\"sv\":\"V2.9.6\",\"os\":\"8.1.1\",\"g\":\"4R30X58T@gmail.com\",\"mid\":\"m828185\",\"nw\":\"3G\",\"l\":\"es\",\"vc\":\"12\",\"hw\":\"640*960\",\"ar\":\"MX\",\"uid\":\"u75770\",\"t\":\"1585901629005\",\"la\":\"6.5\",\"md\":\"HTC-9\",\"vn\":\"1.3.1\",\"ba\":\"HTC\",\"sr\":\"L\"},\"ap\":\"voicebar\",\"et\":[{\"ett\":\"1585970494283\",\"en\":\"start\",\"kv\":{\"entry\":\"3\",\"loading_time\":\"19\",\"action\":\"2\",\"open_ad_type\":\"2\",\"detail\":\"325\"}}]}";
String x = evaluate(line,"mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t");
System.out.println(x);
}
/**
* 解析最初时数据
* */
public static String evaluate(String line, String JsonKeysString){
StringBuilder stringBuilder = new StringBuilder();
//1:先获取所有的key
String[] jsonkeys = JsonKeysString.split(",");
//2:line 将line用|分开来分割,获取时间戳
String[] logContents = line.split("\\|");
//3:校验
if(logContents.length!=2 || StringUtils.isBlank(logContents[1])){
return "";
}
try{
//给json内容创建json对象
JSONObject jsonObject = new JSONObject(logContents[1]);
//获取json内容的公共字段
JSONObject cmjson = jsonObject.getJSONObject("cm");
//开始对json内容的公共字段解析
for(int i=0;i<jsonkeys.length;i++){
String jsonkey = jsonkeys[i].trim();
if(cmjson.has(jsonkey)){
stringBuilder.append(cmjson.getString(jsonkey)).append("\t");
}else {
stringBuilder.append("\t");
}
}
/**
* 开始拼接事件字段和服务器事件咯
* */
stringBuilder.append(jsonObject.getString("et")).append("\t");
stringBuilder.append(logContents[0]).append("\t");
}catch (JSONException e){
e.printStackTrace();
}
return stringBuilder.toString();
}
}
自定义UDTF函数解析非公共“事件”字段
public class etFieldUDTF extends GenericUDTF {
@Override
public StructObjectInspector initialize(StructObjectInspector argOIs) throws UDFArgumentException {
List<String> fieldNames = new ArrayList<>();
List<ObjectInspector> fieldsType = new ArrayList<>();
fieldNames.add("event_name");
fieldsType.add(PrimitiveObjectInspectorFactory.javaStringObjectInspector);
fieldNames.add("event_json");
fieldsType.add(PrimitiveObjectInspectorFactory.javaStringObjectInspector);
return ObjectInspectorFactory.getStandardStructObjectInspector(fieldNames,fieldsType);
}
@Override
public void process(Object[] objects) throws HiveException {
//获取传入的et
String input = objects.toString();
//校验一下
if(StringUtils.isBlank(input)){
return;
}else{
try{
//获取et的Array
JSONArray jsonarrayet = new JSONArray(input);
if (jsonarrayet == null) return;
//遍历每一个事件
for(int i=0;i<jsonarrayet.length();i++){
String[] result = new String[2];
result[0] = jsonarrayet.getJSONObject(i).getString("en");
result[1] = jsonarrayet.getString(i);
forward(result);
}
}catch (JSONException e){
e.printStackTrace();
}
}
}
@Override
public void close() throws HiveException {
}
}
打成jar包,传到hive
运行hive,将jar包添加到运行环境
add jar /usr/local/hive/lib/hiveexecfunction-1.0-SNAPSHOT.jar;

# 创建临时函数
create temporary function base_analizer as 'com.ongbo.UDF.BaseFieldUDF';
create temporary function flat_analizer as 'com.ongbo.UDTF.etFieldUDTF';

## 导入数据到base数据
公共字段
event_name
event_json
server_time
insert overwrite table dwd_base_start_log
partition(dt='2020-04-04')
select
mid_id,user_id,version_code,version_name,lang,source,os,area,model,brand,
sdk_version,gmail,height_width,app_time,network,lng,lat,event_name,event_json,server_time
from
(
select
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[0] as mid_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[1] as user_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[2] as version_code,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[3] as version_name,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[4] as lang,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[5] as source,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[6] as os,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[7] as area,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[8] as model,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[9] as brand,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[10] as sdk_version,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[11] as gmail,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[12] as height_width,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[13] as app_time,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[14] as network,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[15] as lng,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[16] as lat,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[17] as ops,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[18] as server_time
from ods_start_log
where dt='2020-04-04' and base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t')<>''
)
sdk_log lateral view flat_analizer(ops) tmp_k as event_name,event_json;
insert overwrite table dwd_base_event_log
partition(dt='2020-04-04')
select
mid_id,user_id,version_code,version_name,lang,source,os,area,model,brand,
sdk_version,gmail,height_width,app_time,network,lng,lat,event_name,event_json,server_time
from
(
select
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[0] as mid_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[1] as user_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[2] as version_code,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[3] as version_name,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[4] as lang,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[5] as source,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[6] as os,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[7] as area,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[8] as model,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[9] as brand,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[10] as sdk_version,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[11] as gmail,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[12] as height_width,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[13] as app_time,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[14] as network,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[15] as lng,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[16] as lat,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[17] as ops,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[18] as server_time
from ods_event_log
where dt='2020-04-04' and base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t')<>''
)
sdk_log lateral view flat_analizer(ops) tmp_k as event_name,event_json;
上面的只能指定日期,我们可以写成脚本
启动日志导入base数据表的脚本
#~ /bin/bash
APP=voicebar
hive=/usr/local/hive/bin/hive
if [ -n "$1" ] ; then
do_date=$1
else
do_date=`date -d "-1 day" +%F`
fi
sql="
add jar /usr/local/hive/lib/hiveexecfunction-1.0-SNAPSHOT.jar;
create temporary function base_analizer as 'com.ongbo.UDF.BaseFieldUDF';
create temporary function flat_analizer as 'com.ongbo.UDTF.etFieldUDTF';
set hive.exec.dynamic.partition.mode=nonstrict;
insert overwrite table "$APP".dwd_base_start_log
partition(dt='$do_date')
select
mid_id,user_id,version_code,version_name,lang,source,os,area,model,brand,
sdk_version,gmail,height_width,app_time,network,lng,lat,event_name,event_json,server_time
from
(
select
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[0] as mid_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[1] as user_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[2] as version_code,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[3] as version_name,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[4] as lang,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[5] as source,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[6] as os,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[7] as area,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[8] as model,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[9] as brand,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[10] as sdk_version,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[11] as gmail,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[12] as height_width,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[13] as app_time,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[14] as network,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[15] as lng,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[16] as lat,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[17] as ops,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[18] as server_time
from "$APP".ods_start_log
where dt='$do_date' and base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t')<>''
)
sdk_log lateral view flat_analizer(ops) tmp_k as event_name,event_json;
"
$hive -e "$sql"
事件日志的导入base数据表的脚本
#~ /bin/bash
APP=voicebar
hive=/usr/local/hive/bin/hive
if [ -n "$1" ] ; then
do_date=$1
else
do_date=`date -d "-1 day" +%F`
fi
sql="
add jar /usr/local/hive/lib/hiveexecfunction-1.0-SNAPSHOT.jar;
create temporary function base_analizer as 'com.ongbo.UDF.BaseFieldUDF';
create temporary function flat_analizer as 'com.ongbo.UDTF.etFieldUDTF';
set hive.exec.dynamic.partition.mode=nonstrict;
insert overwrite table "$APP".dwd_base_event_log
partition(dt='$do_date')
select
mid_id,user_id,version_code,version_name,lang,source,os,area,model,brand,
sdk_version,gmail,height_width,app_time,network,lng,lat,event_name,event_json,server_time
from
(
select
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[0] as mid_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[1] as user_id,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[2] as version_code,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[3] as version_name,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[4] as lang,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[5] as source,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[6] as os,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[7] as area,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[8] as model,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[9] as brand,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[10] as sdk_version,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[11] as gmail,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[12] as height_width,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[13] as app_time,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[14] as network,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[15] as lng,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[16] as lat,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[17] as ops,
split(base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t'),'\t')[18] as server_time
from "$APP".ods_event_log
where dt='$do_date' and base_analizer(line,'mid,uid,vc,vn,l,sr,os,ar,md,ba,sv,g,hw,nw,ln,la,t')<>''
)
sdk_log lateral view flat_analizer(ops) tmp_k as event_name,event_json;
"
$hive -e "$sql"

到现在,我们的基本数据表dwd_base_start_log和dwd_base_event_log已经完成了,现在就是针对这两张表,不同的时间类型进行展开了。开始吧kk
展开所有的基本表的事件,对事件进行解析,这里的事件包括了启动日志事件和行为日志事件
创建表:
###建表
点击表
create external table dwd_display_log(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`cs` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`action` string,
`goodsid` string,
`place` string,
`extend1` string,
`category` string,
`languagedid` string,
`styleid` string,
`server_time` string)
partitioned by(dt string)
location '/warehouse/voicebar/dwd/dwd_display_log';
详情表
drop table if exists dwd_newsdetail_log;
CREATE EXTERNAL TABLE `dwd_newsdetail_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`entry` string,
`action` string,
`newsid` string,
`languagedid` string,
`news_staytime` string,
`loading_time` string,
`type1` string,
`category` string,
`styleid` string,
`server_time` string)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_newsdetail_log/';
作品列表
drop table if exists dwd_loading_log;
CREATE EXTERNAL TABLE `dwd_loading_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`action` string,
`loading_time` string,
`loading_way` string,
`extend1` string,
`extend2` string,
`type` string,
`type1` string,
`server_time` string)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_loading_log/';
广告表
drop table if exists dwd_ad_log;
CREATE EXTERNAL TABLE `dwd_ad_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`entry` string,
`action` string,
`content` string,
`detail` string,
`ad_source` string,
`behavior` string,
`newstype` string,
`show_style` string,
`server_time` string)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_ad_log/';
消息通知表
drop table if exists dwd_notification_log;
CREATE EXTERNAL TABLE `dwd_notification_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`action` string,
`noti_type` string,
`ap_time` string,
`content` string,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_notification_log/';
用户前台活跃表
drop table if exists dwd_activeforeground_log;
create external table `dwd_activeforeground_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`push_id` string,
`access` string,
`server_time` string)
partitioned by(dt string)
location "/warehouse/voicebar/dwd/dwd_activeforeground_log";
用户后台活跃表
drop table if exists dwd_active_background_log;
CREATE EXTERNAL TABLE `dwd_active_background_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`active_source` string,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_background_log/';
评论表
drop table if exists dwd_comment_log;
CREATE EXTERNAL TABLE `dwd_comment_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`comment_id` int,
`userid` int,
`p_comment_id` int,
`content` string,
`addtime` string,
`other_id` int,
`praise_count` int,
`reply_count` int,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_comment_log/';
收藏表
drop table if exists dwd_favorites_log;
CREATE EXTERNAL TABLE `dwd_favorites_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`id` int,
`course_id` int,
`userid` int,
`add_time` string,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_favorites_log/';
点赞表
drop table if exists dwd_praise_log;
CREATE EXTERNAL TABLE `dwd_praise_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`id` string,
`userid` string,
`target_id` string,
`type` string,
`add_time` string,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_praise_log/';
启动日志表
drop table if exists dwd_start_log;
CREATE EXTERNAL TABLE `dwd_start_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`entry` string,
`open_ad_type` string,
`action` string,
`loading_time` string,
`detail` string,
`extend1` string,
`server_time` string
)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_start_log/';
错误日志表
drop table if exists dwd_error_log;
CREATE EXTERNAL TABLE `dwd_error_log`(
`mid_id` string,
`user_id` string,
`version_code` string,
`version_name` string,
`lang` string,
`source` string,
`os` string,
`area` string,
`model` string,
`brand` string,
`sdk_version` string,
`gmail` string,
`height_width` string,
`app_time` string,
`network` string,
`lng` string,
`lat` string,
`errorBrief` string,
`errorDetail` string,
`server_time` string)
PARTITIONED BY (dt string)
location '/warehouse/voicebar/dwd/dwd_error_log/';
数据导入脚本:
Flink实时数仓
DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。
更多推荐



所有评论(0)