實時
- 訪客主題寬表
- clickhouseUtil
- 商品主題(ProductStatsApp)
- 地區主題ProvinceStatsApp
- 關鍵詞主題KeywordApp
- 資料可視化
- 組件
- 代碼結構
- 需要的依賴
- mapper
- service
- service.impl
- controller
- bean
- keywordStats
- ProductStats
- ProvinceStats
- VisitorStats
- 最終效果
訪客主題寬表

設計一張DWS層的表其實就兩件事:維度和度量(事實資料)
度量包括PV、UV、跳出次數、進入頁面數(session_count)、連續訪問時長
維度包括在分析中比較重要的幾個欄位:渠道、地區、版本、新老用戶進行聚合
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.atguigu.bean.VisitorStats;
import com.atguigu.utils.ClickHouseUtil;
import com.atguigu.utils.MyKafkaUtil;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.api.java.tuple.Tuple4;
import org.apache.flink.streaming.api.datastream.*;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.util.Collector;
import java.text.SimpleDateFormat;
import java.time.Duration;
/**
* mock->nginx->logger->kafka(ods_base_log)->flinkApp(logBaseApp)->
* kafka(dwd_page_log,dwm_unique_visit)->
* flinkApp(UvApp UserJumpApp)->
* kafka(dwm_user_jump_detail)->
* VisitorStatsApp
*/
public class VisitorStatsApp {
public static void main(String[] args) throws Exception {
//1.獲取執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// env.setStateBackend(new FsStateBackend("hdfs://hadoop102:9000/gmall/dwm_log/ck"));
// System.setProperty("HADOOP_USER_NAME", "root");
// env.enableCheckpointing(1000L, CheckpointingMode.EXACTLY_ONCE);
// env.getCheckpointConfig().setCheckpointTimeout(6000L);
//2.讀取kafka主題的資料
// dwd_page_log(pv,訪問時長,進入頁面數)
// dwm_unique_visitor(uv)
// dwm_user_jump(uj)
String groupId = "visitor_stats_app";
String pageViewSourceTopic = "dwd_page_log";
String uniqueVisitSourceTopic = "dwm_unique_visit";
String userJumpDetailSourceTopic = "dwm_user_jump_detail";
DataStreamSource<String> pageLogDS = env.addSource(MyKafkaUtil.getKafkaSource(pageViewSourceTopic, groupId));
DataStreamSource<String> uvDS = env.addSource(MyKafkaUtil.getKafkaSource(uniqueVisitSourceTopic, groupId));
DataStreamSource<String> userJumpDS = env.addSource(MyKafkaUtil.getKafkaSource(userJumpDetailSourceTopic, groupId));
//測驗
// pageLogDS.print("page>>>>>>>>>>");
// uvDS.print("uv>>>>>>>");
// userJumpDS.print("userJump>>>>>>>");
//3.格式化流資料,使其欄位統一(JavaBean)
//3.1將頁面資料流格式化為visitorStats,PV during-time
SingleOutputStreamOperator<VisitorStats> pvAndDtDS = pageLogDS.map(jsonStr -> {
//將資料轉換成JSON物件
JSONObject jsonObject = JSON.parseObject(jsonStr);
JSONObject commonObj = jsonObject.getJSONObject("common");
return new VisitorStats("", "",
commonObj.getString("vc"),
commonObj.getString("ch"),
commonObj.getString("ar"),
commonObj.getString("is_new"),
0L,
1L,
0L,
0L,
jsonObject.getJSONObject("page").getLong("during_time"),
jsonObject.getLong("ts"));
});
//3.2 將頁面資料流先過濾后格式化為visitorStats sv_ct
SingleOutputStreamOperator<VisitorStats> svCountDS = pageLogDS.process(new ProcessFunction<String, VisitorStats>() {
@Override
public void processElement(String value, Context ctx, Collector<VisitorStats> out) throws Exception {
//將資料轉換成JSON物件
JSONObject jsonObject = JSON.parseObject(value);
//獲取上一條頁面資料
String lastPage = jsonObject.getJSONObject("page").getString("last_page_id");
if (lastPage == null || lastPage.length() <= 0) {
JSONObject commonObj = jsonObject.getJSONObject("common");
out.collect(new VisitorStats("", "",
commonObj.getString("vc"),
commonObj.getString("ch"),
commonObj.getString("ar"),
commonObj.getString("is_new"),
0L,
0L,
1L,
0L,
0L,
jsonObject.getLong("ts"))
);
}
}
});
//3.3 將uvDS格式化為VisitorStats uv
SingleOutputStreamOperator<VisitorStats> uvCountDS = uvDS.map(jsonStr -> {
JSONObject jsonObject = JSON.parseObject(jsonStr);
JSONObject commonObj = jsonObject.getJSONObject("common");
return new VisitorStats("", "",
commonObj.getString("vc"),
commonObj.getString("ch"),
commonObj.getString("ar"),
commonObj.getString("is_new"),
1L,
0L,
0L,
0L,
0L,
jsonObject.getLong("ts")
);
});
//3.4 將userJumpDS格式化為VisitorStats uj_Ct
SingleOutputStreamOperator<VisitorStats> ujCountDS = userJumpDS.map(jsonStr -> {
JSONObject jsonObject = JSON.parseObject(jsonStr);
JSONObject commonObj = jsonObject.getJSONObject("common");
return new VisitorStats("", "",
commonObj.getString("vc"),
commonObj.getString("ch"),
commonObj.getString("ar"),
commonObj.getString("is_new"),
0L,
0L,
0L,
1L,
0L,
jsonObject.getLong("ts"));
});
//4.將多個流的資料進行union
DataStream<VisitorStats> unionDS = pvAndDtDS.union(svCountDS, uvCountDS, ujCountDS);
SingleOutputStreamOperator<VisitorStats> visitorStatsSingleOutputStreamOperator = unionDS.assignTimestampsAndWatermarks(WatermarkStrategy
.<VisitorStats>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner(new SerializableTimestampAssigner<VisitorStats>() {
@Override
public long extractTimestamp(VisitorStats element, long recordTimestamp) {
return element.getTs();
}
}));
//5.分組,聚合計算
KeyedStream<VisitorStats, Tuple4<String, String, String, String>> keyedStream = visitorStatsSingleOutputStreamOperator
.keyBy(new KeySelector<VisitorStats, Tuple4<String, String, String, String>>() {
@Override
public Tuple4<String, String, String, String> getKey(VisitorStats value) throws Exception {
return Tuple4.of(value.getVc(), value.getCh(), value.getAr(), value.getIs_new());
}
});
//開窗
WindowedStream<VisitorStats, Tuple4<String, String, String, String>, TimeWindow> windowedStream =
keyedStream.window(TumblingEventTimeWindows.of(Time.seconds(10)));
//聚合操作
SingleOutputStreamOperator<Object> result = windowedStream.reduce(
new ReduceFunction<VisitorStats>() {
@Override
public VisitorStats reduce(VisitorStats value1, VisitorStats value2) throws Exception {
return new VisitorStats("", "",
value1.getVc(),
value1.getCh(),
value1.getAr(),
value1.getIs_new(),
value1.getUv_ct() + value2.getUv_ct(),
value1.getPv_ct() + value2.getPv_ct(),
value1.getSv_ct() + value2.getSv_ct(),
value1.getUj_ct() + value2.getUj_ct(),
value1.getDur_sum() + value2.getDur_sum(),
System.currentTimeMillis());
}
},
new WindowFunction<VisitorStats, Object, Tuple4<String, String, String, String>, TimeWindow>() {
@Override
public void apply(Tuple4<String, String, String, String> stringStringStringStringTuple4, TimeWindow window, Iterable<VisitorStats> input, Collector<Object> out) throws Exception {
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
//取出資料
VisitorStats visitorStats = input.iterator().next();
//取出視窗的開始及結束時間
long start = window.getStart();
long end = window.getEnd();
String stt = sdf.format(start);
String edt = sdf.format(end);
//設定時間
visitorStats.setStt(stt);
visitorStats.setEdt(edt);
//將資料寫出
out.collect(visitorStats);
}
});
result.print(">>>>>>>>>");
//6.將聚合之后的資料寫入clickHouse
result.addSink(ClickHouseUtil.getSink("insert into visitor_stats values(?,?,?,?,?,?,?,?,?,?,?,?)"));
//7.執行任務
env.execute();
}
}
clickhouseUtil
import com.atguigu.bean.TransientSink;
import com.atguigu.common.GmallConfig;
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
import java.beans.Transient;
import java.lang.reflect.Field;
import java.sql.PreparedStatement;
import java.sql.SQLException;
public class ClickHouseUtil {
public static <T> SinkFunction getSink(String sql) {
return JdbcSink.sink(sql, new JdbcStatementBuilder<T>() {
@Override
public void accept(PreparedStatement preparedStatement, T obj) throws SQLException {
//反射的方式獲取所有的屬性名
Field[] fields = obj.getClass().getDeclaredFields();
//定義跳過的屬性
int offset = 0;
for (int i = 0; i < fields.length; i++) {
//獲取欄位名
Field field = fields[i];
//獲取欄位上的注解
TransientSink transientSink = field.getAnnotation(TransientSink.class);
if (transientSink != null) {
offset++;
continue;
}
//設定可訪問私有屬性的值
field.setAccessible(true);
try {
Object o = field.get(obj);
//給占位符賦值
preparedStatement.setObject(i + 1 - offset, o);
} catch (IllegalAccessException e) {
e.printStackTrace();
}
}
}
},
JdbcExecutionOptions.builder()
.withBatchSize(5)
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withDriverName(GmallConfig.CLICKHOUSE_DRIVER)
.withUrl(GmallConfig.CLICKHOUSE_URL)
.build());
}
}
商品主題(ProductStatsApp)

package com.atguigu.app.dws;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.atguigu.app.func.DimAsyncFunction;
import com.atguigu.bean.OrderWide;
import com.atguigu.bean.PaymentWide;
import com.atguigu.bean.ProductStats;
import com.atguigu.utils.ClickHouseUtil;
import com.atguigu.utils.DateTimeUtil;
import com.atguigu.utils.GmallConstant;
import com.atguigu.utils.MyKafkaUtil;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.streaming.api.datastream.AsyncDataStream;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.util.Collector;
import java.math.BigDecimal;
import java.text.SimpleDateFormat;
import java.time.Duration;
import java.util.Collections;
import java.util.HashSet;
import java.util.concurrent.TimeUnit;
/**
* mockLog -> Nginx -> Logger -> Kafka(ods_base_log) -> FlinkApp(BaseLogApp) -> Kafka(dwd_page_log)
* mockDB -> MySQL -> Maxwell -> Kafka(ods_base_db_m) -> FlinkApp(BaseDBApp) -> Kafka(HBase) ->
* FlinkApp(OrderWideApp,PaymentWideApp,Redis,Phoenix) -> Kafka(dwm_order_wide,dwm_payment_wide)
*/
public class ProductStatsApp {
public static void main(String[] args) throws Exception {
//1.獲取執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
//1.1 設定狀態后端
//env.setStateBackend(new FsStateBackend("hdfs://hadoop102:8020/gmall/dwd_log/ck"));
//1.2 開啟CK
//env.enableCheckpointing(10000L, CheckpointingMode.EXACTLY_ONCE);
//env.getCheckpointConfig().setCheckpointTimeout(60000L);
//2.讀取Kafka資料 7個主題
String groupId = "product_stats_app1";
String pageViewSourceTopic = "dwd_page_log";
String favorInfoSourceTopic = "dwd_favor_info";
String cartInfoSourceTopic = "dwd_cart_info";
String orderWideSourceTopic = "dwm_order_wide";
String paymentWideSourceTopic = "dwm_payment_wide";
String refundInfoSourceTopic = "dwd_order_refund_info";
String commentInfoSourceTopic = "dwd_comment_info";
FlinkKafkaConsumer<String> pageViewSource = MyKafkaUtil.getKafkaSource(pageViewSourceTopic, groupId);
DataStreamSource<String> pageViewDStream = env.addSource(pageViewSource);
FlinkKafkaConsumer<String> favorInfoSourceSource = MyKafkaUtil.getKafkaSource(favorInfoSourceTopic, groupId);
DataStreamSource<String> favorInfoDStream = env.addSource(favorInfoSourceSource);
FlinkKafkaConsumer<String> cartInfoSource = MyKafkaUtil.getKafkaSource(cartInfoSourceTopic, groupId);
DataStreamSource<String> cartInfoDStream = env.addSource(cartInfoSource);
FlinkKafkaConsumer<String> orderWideSource = MyKafkaUtil.getKafkaSource(orderWideSourceTopic, groupId);
DataStreamSource<String> orderWideDStream = env.addSource(orderWideSource);
FlinkKafkaConsumer<String> paymentWideSource = MyKafkaUtil.getKafkaSource(paymentWideSourceTopic, groupId);
DataStreamSource<String> paymentWideDStream = env.addSource(paymentWideSource);
FlinkKafkaConsumer<String> refundInfoSource = MyKafkaUtil.getKafkaSource(refundInfoSourceTopic, groupId);
DataStreamSource<String> refundInfoDStream = env.addSource(refundInfoSource);
FlinkKafkaConsumer<String> commentInfoSource = MyKafkaUtil.getKafkaSource(commentInfoSourceTopic, groupId);
DataStreamSource<String> commentInfoDStream = env.addSource(commentInfoSource);
//3.將7個流轉換為統一資料格式
//3.1 處理點擊和曝光資料
SingleOutputStreamOperator<ProductStats> clickAndDisplayDS = pageViewDStream.process(new ProcessFunction<String, ProductStats>() {
@Override
public void processElement(String pageLog, Context ctx, Collector<ProductStats> out) throws Exception {
//轉換為JSON物件
JSONObject jsonObject = JSON.parseObject(pageLog);
//獲取PageId
JSONObject pageObj = jsonObject.getJSONObject("page");
String page_id = pageObj.getString("page_id");
//獲取時間戳欄位
Long ts = jsonObject.getLong("ts");
//判斷是good_detail頁面,則為點擊資料
if ("good_detail".equals(page_id)) {
ProductStats productStats = ProductStats
.builder()
.sku_id(pageObj.getLong("item"))
.click_ct(1L)
.ts(ts)
.build();
out.collect(productStats);
}
//取出曝光資料
JSONArray displays = jsonObject.getJSONArray("displays");
if (displays != null && displays.size() > 0) {
for (int i = 0; i < displays.size(); i++) {
JSONObject displayObj = displays.getJSONObject(i);
if ("sku_id".equals(displayObj.getString("item_type"))) {
ProductStats productStats = ProductStats.builder()
.sku_id(displayObj.getLong("item"))
.display_ct(1L)
.ts(ts)
.build();
out.collect(productStats);
}
}
}
}
});
//3.2 處理收藏資料
SingleOutputStreamOperator<ProductStats> favorStatsDS = favorInfoDStream.map(
json -> {
JSONObject favorInfo = JSON.parseObject(json);
Long ts = DateTimeUtil.toTs(favorInfo.getString("create_time"));
return ProductStats.builder()
.sku_id(favorInfo.getLong("sku_id"))
.favor_ct(1L)
.ts(ts)
.build();
});
//3.3 處理加購資料
SingleOutputStreamOperator<ProductStats> cartStatsDS = cartInfoDStream.map(
json -> {
JSONObject cartInfo = JSON.parseObject(json);
Long ts = DateTimeUtil.toTs(cartInfo.getString("create_time"));
return ProductStats.builder()
.sku_id(cartInfo.getLong("sku_id"))
.cart_ct(1L)
.ts(ts)
.build();
});
//3.4 處理下單資料
SingleOutputStreamOperator<ProductStats> orderDS = orderWideDStream.map(new MapFunction<String, ProductStats>() {
@Override
public ProductStats map(String value) throws Exception {
//轉換為OrderWide物件
OrderWide orderWide = JSON.parseObject(value, OrderWide.class);
//獲取時間戳欄位
Long ts = DateTimeUtil.toTs(orderWide.getCreate_time());
return ProductStats.builder()
.sku_id(orderWide.getSku_id())
.order_amount(orderWide.getTotal_amount())
.order_sku_num(orderWide.getSku_num())
.orderIdSet(new HashSet<>(Collections.singleton(orderWide.getOrder_id())))
.ts(ts)
.build();
}
});
//3.5 處理支付資料
SingleOutputStreamOperator<ProductStats> paymentStatsDS = paymentWideDStream.map(
json -> {
PaymentWide paymentWide = JSON.parseObject(json, PaymentWide.class);
Long ts = DateTimeUtil.toTs(paymentWide.getPayment_create_time());
return ProductStats.builder()
.sku_id(paymentWide.getSku_id())
.payment_amount(paymentWide.getSplit_total_amount())
.paidOrderIdSet(new HashSet<>(Collections.singleton(paymentWide.getOrder_id())))
.ts(ts).build();
});
//3.6 處理退單資料
SingleOutputStreamOperator<ProductStats> refundStatsDS = refundInfoDStream.map(
json -> {
JSONObject refundJsonObj = JSON.parseObject(json);
Long ts = DateTimeUtil.toTs(refundJsonObj.getString("create_time"));
return ProductStats.builder()
.sku_id(refundJsonObj.getLong("sku_id"))
.refund_amount(refundJsonObj.getBigDecimal("refund_amount"))
.refundOrderIdSet(new HashSet<>(Collections.singleton(refundJsonObj.getLong("order_id"))))
.ts(ts)
.build();
});
//3.7 處理評價資料
SingleOutputStreamOperator<ProductStats> appraiseDS = commentInfoDStream.map(new MapFunction<String, ProductStats>() {
@Override
public ProductStats map(String value) throws Exception {
//將資料轉換為JSON物件
JSONObject jsonObject = JSON.parseObject(value);
Long ts = DateTimeUtil.toTs(jsonObject.getString("create_time"));
//處理好評數
Long goodCommentCt = 0L;
if (GmallConstant.APPRAISE_GOOD.equals(jsonObject.getString("appraise"))) {
goodCommentCt = 1L;
}
return ProductStats.builder()
.sku_id(jsonObject.getLong("sku_id"))
.comment_ct(1L)
.good_comment_ct(goodCommentCt)
.ts(ts)
.build();
}
});
//4.Union
DataStream<ProductStats> unionDS = clickAndDisplayDS.union(favorStatsDS,
cartStatsDS,
orderDS,
paymentStatsDS,
refundStatsDS,
appraiseDS);
//5.設定Watermark
SingleOutputStreamOperator<ProductStats> productStatsWithWaterMarkDS = unionDS.assignTimestampsAndWatermarks(WatermarkStrategy.<ProductStats>forBoundedOutOfOrderness(Duration.ofSeconds(10L)).withTimestampAssigner(new SerializableTimestampAssigner<ProductStats>() {
@Override
public long extractTimestamp(ProductStats element, long recordTimestamp) {
return element.getTs();
}
}));
//6.分組、開窗、聚合
SingleOutputStreamOperator<ProductStats> reduceDS = productStatsWithWaterMarkDS.keyBy(ProductStats::getSku_id)
.window(TumblingEventTimeWindows.of(Time.seconds(2)))
.reduce(new ReduceFunction<ProductStats>() {
@Override
public ProductStats reduce(ProductStats stats1, ProductStats stats2) throws Exception {
stats1.setDisplay_ct(stats1.getDisplay_ct() + stats2.getDisplay_ct());
stats1.setClick_ct(stats1.getClick_ct() + stats2.getClick_ct());
stats1.setCart_ct(stats1.getCart_ct() + stats2.getCart_ct());
stats1.setFavor_ct(stats1.getFavor_ct() + stats2.getFavor_ct());
stats1.setOrder_sku_num(stats1.getOrder_sku_num() + stats2.getOrder_sku_num());
stats1.setOrder_amount(stats1.getOrder_amount().add(stats2.getOrder_amount()));
stats1.getOrderIdSet().addAll(stats2.getOrderIdSet());
stats1.setOrder_ct((long) stats1.getOrderIdSet().size());
stats1.setPayment_amount(stats1.getPayment_amount().add(stats2.getPayment_amount()));
stats1.getPaidOrderIdSet().addAll(stats2.getPaidOrderIdSet());
stats1.getRefundOrderIdSet().addAll(stats2.getRefundOrderIdSet());
stats1.setRefund_order_ct((long) stats1.getRefundOrderIdSet().size());
stats1.setRefund_amount(stats1.getRefund_amount().add(stats2.getRefund_amount()));
stats1.setPaid_order_ct((long) stats1.getPaidOrderIdSet().size());
stats1.setComment_ct(stats1.getComment_ct() + stats2.getComment_ct());
stats1.setGood_comment_ct(stats1.getGood_comment_ct() + stats2.getGood_comment_ct());
return stats1;
}
}, new WindowFunction<ProductStats, ProductStats, Long, TimeWindow>() {
@Override
public void apply(Long aLong, TimeWindow window, Iterable<ProductStats> input, Collector<ProductStats> out) throws Exception {
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
long start = window.getStart();
long end = window.getEnd();
String stt = sdf.format(start);
String edt = sdf.format(end);
//取出聚合以后的資料
ProductStats productStats = input.iterator().next();
//設定視窗時間
productStats.setStt(stt);
productStats.setEdt(edt);
out.collect(productStats);
}
});
//7.關聯維度資訊
//7.1 關聯SKU資訊
SingleOutputStreamOperator<ProductStats> productStatsWithSkuDS = AsyncDataStream.unorderedWait(reduceDS, new DimAsyncFunction<ProductStats>("DIM_SKU_INFO") {
@Override
public String getKey(ProductStats productStats) {
return productStats.getSku_id().toString();
}
@Override
public void join(ProductStats productStats, JSONObject dimInfo) throws Exception {
//獲取維度中的資訊
String sku_name = dimInfo.getString("SKU_NAME");
BigDecimal price = dimInfo.getBigDecimal("PRICE");
Long spu_id = dimInfo.getLong("SPU_ID");
Long tm_id = dimInfo.getLong("TM_ID");
Long category3_id = dimInfo.getLong("CATEGORY3_ID");
//關聯SKU維度資訊
productStats.setSku_name(sku_name);
productStats.setSku_price(price);
productStats.setSpu_id(spu_id);
productStats.setTm_id(tm_id);
productStats.setCategory3_id(category3_id);
}
}, 300, TimeUnit.SECONDS);
//7.2 補充SPU維度
SingleOutputStreamOperator<ProductStats> productStatsWithSpuDS =
AsyncDataStream.unorderedWait(productStatsWithSkuDS,
new DimAsyncFunction<ProductStats>("DIM_SPU_INFO") {
@Override
public void join(ProductStats productStats, JSONObject jsonObject) throws Exception {
productStats.setSpu_name(jsonObject.getString("SPU_NAME"));
}
@Override
public String getKey(ProductStats productStats) {
return String.valueOf(productStats.getSpu_id());
}
}, 60, TimeUnit.SECONDS);
//7.3 補充品類維度
SingleOutputStreamOperator<ProductStats> productStatsWithCategory3DS =
AsyncDataStream.unorderedWait(productStatsWithSpuDS,
new DimAsyncFunction<ProductStats>("DIM_BASE_CATEGORY3") {
@Override
public void join(ProductStats productStats, JSONObject jsonObject) throws Exception {
productStats.setCategory3_name(jsonObject.getString("NAME"));
}
@Override
public String getKey(ProductStats productStats) {
return String.valueOf(productStats.getCategory3_id());
}
}, 60, TimeUnit.SECONDS);
//7.4 補充品牌維度
SingleOutputStreamOperator<ProductStats> productStatsWithTmDS =
AsyncDataStream.unorderedWait(productStatsWithCategory3DS,
new DimAsyncFunction<ProductStats>("DIM_BASE_TRADEMARK") {
@Override
public void join(ProductStats productStats, JSONObject jsonObject) throws Exception {
productStats.setTm_name(jsonObject.getString("TM_NAME"));
}
@Override
public String getKey(ProductStats productStats) {
return String.valueOf(productStats.getTm_id());
}
}, 60, TimeUnit.SECONDS);
//列印測驗
productStatsWithTmDS.print();
//8.寫入ClickHouse
productStatsWithTmDS.addSink(ClickHouseUtil.<ProductStats>getSink("insert into product_stats_200821 values(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)"));
//9.執行任務
env.execute();
}
}
地區主題ProvinceStatsApp
package com.atguigu.app.dws;
import com.atguigu.bean.ProvinceStats;
import com.atguigu.utils.ClickHouseUtil;
import com.atguigu.utils.MyKafkaUtil;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
public class ProvinceStatsApp {
public static void main(String[] args) throws Exception {
//1.獲取執行環境(流、表)
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// env.setStateBackend(new FsStateBackend("hdfs://hadoop102:9000/gmall/dwm_log/ck"));
// System.setProperty("HADOOP_USER_NAME", "root");
// env.enableCheckpointing(1000L, CheckpointingMode.EXACTLY_ONCE);
// env.getCheckpointConfig().setCheckpointTimeout(6000L);
//獲取表的執行環境
EnvironmentSettings settings = EnvironmentSettings
.newInstance()
.inStreamingMode()
.useBlinkPlanner()
.build();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings);
//2.讀取kafka資料創建動態表
String sourceTopic = "dwm_order_wide";
String groupId = "province_stats_app";
tableEnv.executeSql("CREATE TABLE ORDER_WIDE (province_id BIGINT, " +
"province_name STRING," +
"province_area_code STRING," +
"province_iso_code STRING," +
"province_3166_2_code STRING," +
"order_id STRING, " +
"split_total_amount DOUBLE," +
"create_time STRING," +
"rowtime AS TO_TIMESTAMP(create_time)," +
"WATERMARK FOR rowtime AS rowtime)" +
" WITH (" + MyKafkaUtil.getKafkaDDL(sourceTopic, groupId) + ")");
//3.分組、開窗、聚合
Table reduceTable = tableEnv.sqlQuery("select " +
"TUMBLE_START(rowtime,INTERVAL '10' SECOND) as stt,"+
"TUMBLE_END(rowtime,INTERVAL '10' SECOND) as edt"+
"province_id," +
"province_name," +
"province_area_code," +
"province_iso_code," +
"province_3166_2_code ," +
"sum(split_total_amount) order_amount," +
"count(*)order_count,"+
"UNIX_TIMESTAMP()*1000 ts " +
"from ORDER_WIDE " +
"group by province_id,province_name,province_area_code,province_iso_code,province_3166_2_code,TUMBLE(rowtime INTERVAL '10' SECOND)");
//4.將動態表轉為追加流
DataStream<ProvinceStats> rowDataStream = tableEnv.toAppendStream(reduceTable, ProvinceStats.class);
rowDataStream.print();
//5.寫入clickhouse
rowDataStream.addSink(ClickHouseUtil.<ProvinceStats>getSink("insert into province_stats_2021 values(?,?,?,?,?,?,?,?,?,?)"));
//6.執行任務
env.execute();
}
}
關鍵詞主題KeywordApp
從Kafka主題中獲得資料流
把Json字串資料流轉換為統一資料物件的資料流
把統一的資料結構流合并為一個流
設定事件時間與水位線
分組、開窗、聚合
寫入ClickHousea
package com.atguigu.app.dws;
import com.atguigu.app.func.KeyWordUDTF;
import com.atguigu.bean.KeywordStats;
import com.atguigu.utils.ClickHouseUtil;
import com.atguigu.utils.GmallConstant;
import com.atguigu.utils.KeyWordUtil;
import com.atguigu.utils.MyKafkaUtil;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.reflections.vfs.SystemFile;
public class KeyWordStatsApp {
private static SystemFile systemFile;
public static void main(String[] args) throws Exception {
//1、獲取執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// env.setStateBackend(new FsStateBackend("hdfs://hadoop102:9000/gmall/dwm_log/ck"));
// System.setProperty("HADOOP_USER_NAME", "root");
// env.enableCheckpointing(1000L, CheckpointingMode.EXACTLY_ONCE);
// env.getCheckpointConfig().setCheckpointTimeout(6000L);
//獲取表的執行環境
EnvironmentSettings settings = EnvironmentSettings
.newInstance()
.inStreamingMode()
.useBlinkPlanner()
.build();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings);
//2、讀取kafka主題資料創建動態表 dwd_page_log
String sourceTopic = "dwd_page_log";
String groupId = "keyword_stats_app";
tableEnv.executeSql("create table page_view(" +
" common MAP<STRING,STRING>," +
" page MAP<STRING,STRING>," +
" ts BIGINT," +
" rowtime AS TO_TIMESTAMP(FROM_UNIXTIME(ts/1000,'yyyy-MM-dd HH:mm:ss'))," +
" WATERMARK FOR rowtime AS rowtime - INTERVAL '2' SECOND)" +
" WITH (" + MyKafkaUtil.getKafkaDDL(sourceTopic, groupId) + ")");
//3、過濾資料,只需要搜索的資料,搜索的關鍵詞不能為空
Table filterTable = tableEnv.sqlQuery("select page['item'] fullWord,rowtime" +
" from page_view" +
" where page['item_type']='keyword' and page['item'] IS NOT NULL");
//4、使用udtf函式進行切詞處理 函式注冊
tableEnv.createTemporarySystemFunction("ik_analyze", KeyWordUDTF.class);
Table analyzeWordTable = tableEnv.sqlQuery("SELECT" +
" keyword," +
" rowtime " +
" from " + filterTable +
" ,LATERAL TABLE(ik_analyze(fullWord)) AS T(keyword)");
//5、分組、開窗、聚合
Table resultTable = tableEnv.sqlQuery("select keyword," +
"count(*) ct, '"
+ GmallConstant.KEYWORD_SEARCH + "' source ," +
"DATE_FORMAT(TUMBLE_START(rowtime, INTERVAL '10' SECOND),'yyyy-MM-dd HH:mm:ss') stt," +
"DATE_FORMAT(TUMBLE_END(rowtime, INTERVAL '10' SECOND),'yyyy-MM-dd HH:mm:ss') edt," +
"UNIX_TIMESTAMP()*1000 ts from " + analyzeWordTable
+ " GROUP BY TUMBLE(rowtime, INTERVAL '10' SECOND ),keyword");
DataStream<KeywordStats> keywordStatsDataStream = tableEnv.toAppendStream(resultTable, KeywordStats.class);
keywordStatsDataStream.print();
//6、將資料寫入clickhouse
keywordStatsDataStream.addSink(
ClickHouseUtil.<KeywordStats>getSink(
"insert into keyword_stats_200821(keyword,ct,source,stt,edt,ts) " +
" values(?,?,?,?,?,?)")
);
//7、執行任務
env.execute();
}
}
資料可視化
組件
| 組件名稱 | 組件 | 查詢指標 | 對應的資料表 |
|---|---|---|---|
| 總成交金額 | 數字翻牌 | 訂單總金額 | product_stats |
| 省市熱力圖查詢 | 熱力圖 | 省市分組訂單金額 | province_stats |
| 分時流量 | 折線圖 | UV分時數 PV分時數 新用戶分時數 | visitor_stats |
| 品牌TopN | 水平柱狀圖 | 按品牌分組訂單金額 | product_stats |
| 品類分布 | 餅狀圖 | 按品類分組訂單金額 | product_stats |
| 熱詞字符云 | 字符云 | 關鍵詞分組計數 | keyword_stats |
| 流量表格 | 交叉透視表 | UV數(新老用戶) PV數(新老用戶) 跳出率(新老用戶) 平均訪問時長 (新老用戶) 平均訪問頁面數(新老用戶) | visitor_stats |
| 熱門商品 | 輪播表格 | 按SPU分組訂單金額 | product_stats |
代碼結構
| 分層 | 類 | 處理內容 |
|---|---|---|
| controller 控制層 | SugarController | 查詢交易額介面及回傳引數處理 |
| service 服務層 | ProductStatsService ProductStatsServiceImpl | 查詢商品統計資料 |
| mapper 資料映射層 | ProductStatsMapper | 撰寫SQL查詢商品統計表 |
需要的依賴
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.4.2</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<groupId>com.atguigu.gmall</groupId>
<artifactId>gmall2021-publisher</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>gmall2021-publisher</name>
<description>Demo project for Spring Boot</description>
<properties>
<java.version>1.8</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>2.1.3</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
<exclusions>
<exclusion>
<groupId>org.junit.vintage</groupId>
<artifactId>junit-vintage-engine</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.11</version>
</dependency>
<dependency>
<groupId>ru.yandex.clickhouse</groupId>
<artifactId>clickhouse-jdbc</artifactId>
<version>0.1.55</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
mapper
注代碼都放一起了,沒有細分
package com.example.gmallpublisher.mapper;
import com.example.gmallpublisher.bean.ProductStats;
import com.example.gmallpublisher.bean.ProvinceStats;
import com.example.gmallpublisher.bean.VisitorStats;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
import java.math.BigDecimal;
import java.util.List;
public interface ProductMapper {
//1、查詢GMV總數
@Select("select sum(order_amount) order_amount from product_stats_200821 where toYYYYMMDD(stt)=#{date}")
BigDecimal getSumAmount(int date);
//統計某天不同SPU商品交易額排名
@Select("select spu_id,spu_name,sum(order_amount) order_amount," +
"sum(order_ct) order_ct from product_stats_200821 " +
"where toYYYYMMDD(stt)=#{date} group by spu_id,spu_name " +
"having order_amount>0 order by order_amount desc limit #{limit} ")
List<ProductStats> getProductStatsGroupBySpu(@Param("date") int date, @Param("limit") int limit);
//統計某天不同類別商品交易額排名
@Select("select category3_id,category3_name,sum(order_amount) order_amount " +
"from product_stats_200821 " +
"where toYYYYMMDD(stt)=#{date} group by category3_id,category3_name " +
"having order_amount>0 order by order_amount desc limit #{limit}")
List<ProductStats> getProductStatsGroupByCategory3(@Param("date")int date , @Param("limit") int limit);
//統計某天不同品牌商品交易額排名
@Select("select tm_id,tm_name,sum(order_amount) order_amount " +
"from product_stats_200821 " +
"where toYYYYMMDD(stt)=#{date} group by tm_id,tm_name " +
"having order_amount=0 order by order_amount desc limit #{limit} ")
List<ProductStats> getProductStatsByTrademark(@Param("date")int date, @Param("limit") int limit);
//按地區查詢交易額
@Select("select province_name,sum(order_amount) order_amount " +
"from province_stats_200821 where toYYYYMMDD(stt)=#{date} " +
"group by province_id ,province_name")
List<ProvinceStats> selectProvinceStats(int date);
//新老訪客流量統計
@Select("select is_new,sum(uv_ct) uv_ct,sum(pv_ct) pv_ct," +
"sum(sv_ct) sv_ct, sum(uj_ct) uj_ct,sum(dur_sum) dur_sum " +
"from visitor_stats where toYYYYMMDD(stt)=#{date} group by is_new")
List<VisitorStats> selectVisitorStatsByNewFlag(int date);
//分時流量統計
@Select("select sum(if(is_new='1', visitor_stats.uv_ct,0)) new_uv,toHour(stt) hr," +
"sum(visitor_stats.uv_ct) uv_ct, sum(pv_ct) pv_ct, sum(uj_ct) uj_ct " +
"from visitor_stats where toYYYYMMDD(stt)=#{date} group by toHour(stt)")
List<VisitorStats> selectVisitorStatsByHour(int date);
@Select("select count(pv_ct) pv_ct from visitor_stats " +
"where toYYYYMMDD(stt)=#{date} ")
Long selectPv(int date);
@Select("select count(uv_ct) uv_ct from visitor_stats " +
"where toYYYYMMDD(stt)=#{date} ")
Long selectUv(int date);
}
service
package com.example.gmallpublisher.service;
import com.example.gmallpublisher.bean.ProductStats;
import com.example.gmallpublisher.bean.ProvinceStats;
import com.example.gmallpublisher.bean.VisitorStats;
import javax.websocket.RemoteEndpoint;
import java.math.BigDecimal;
import java.util.List;
public interface ProductService {
BigDecimal getSumAmount(int date);
//統計某天不同SPU商品交易額排名
List<ProductStats> getProductStatsGroupBySpu(int date, int limit);
//統計某天不同類別商品交易額排名
List<ProductStats> getProductStatsGroupByCategory3(int date,int limit);
//統計某天不同品牌商品交易額排名
List<ProductStats> getProductStatsByTrademark(int date,int limit);
/**
* Desc: 地區維度統計介面
*/
List<ProvinceStats> getProvinceStats(int date);
List<VisitorStats> getVisitorStatsByNewFlag(int date);
List<VisitorStats> getVisitorStatsByHour(int date);
Long getPv(int date);
Long getUv(int date);
}
service.impl
package com.example.gmallpublisher.service.impl;
import com.example.gmallpublisher.bean.ProductStats;
import com.example.gmallpublisher.bean.ProvinceStats;
import com.example.gmallpublisher.bean.VisitorStats;
import com.example.gmallpublisher.mapper.ProductMapper;
import com.example.gmallpublisher.service.ProductService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.util.List;
@Service
public class ProductServiceImpl implements ProductService {
@Autowired
ProductMapper productMapper;
@Override
public BigDecimal getSumAmount(int date) {
return productMapper.getSumAmount(date);
}
@Override
public List<ProductStats> getProductStatsGroupBySpu(int date, int limit) {
return productMapper.getProductStatsGroupBySpu(date, limit);
}
@Override
public List<ProductStats> getProductStatsGroupByCategory3(int date, int limit) {
return productMapper.getProductStatsGroupByCategory3(date, limit);
}
@Override
public List<ProductStats> getProductStatsByTrademark(int date, int limit) {
return productMapper.getProductStatsByTrademark(date, limit);
}
@Override
public List<ProvinceStats> getProvinceStats(int date) {
return productMapper.selectProvinceStats(date);
}
@Override
public List<VisitorStats> getVisitorStatsByNewFlag(int date) {
return productMapper.selectVisitorStatsByNewFlag(date);
}
@Override
public List<VisitorStats> getVisitorStatsByHour(int date) {
return productMapper.selectVisitorStatsByHour(date);
}
@Override
public Long getPv(int date) {
return productMapper.selectPv(date);
}
@Override
public Long getUv(int date) {
return productMapper.selectUv(date);
}
}
controller
package com.example.gmallpublisher.controller;
import com.example.gmallpublisher.bean.KeywordStats;
import com.example.gmallpublisher.bean.ProductStats;
import com.example.gmallpublisher.bean.ProvinceStats;
import com.example.gmallpublisher.bean.VisitorStats;
import com.example.gmallpublisher.service.KeywordService;
import com.example.gmallpublisher.service.ProductService;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.math.BigDecimal;
import java.text.SimpleDateFormat;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
//@Controller 回傳頁面
@RestController // =@Controller + @ResponseBody
@RequestMapping("/api/sugar")
public class SugarController {
//@RequestMapping("/gmv")
/**
* {
* "stats":0,
* "msg":"",
* "data":1201083.4767158988
* }
*/
@Autowired
ProductService productService;
@RequestMapping("/gmv")
public String getGmv(@RequestParam(value = "date", defaultValue = "0") int date) {
if (date == 0) {
date = now();
}
return "{" +
"\"status\": 0," +
"\"msg\":\"\"," +
"\"data\":" + productService.getSumAmount(date) +
"}";
}
private int now() {
SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMdd");
return Integer.parseInt(sdf.format(new Date().getTime()));
}
@RequestMapping("/test")
public String test1() {
return "hello";
}
//商品串列介面方法
@RequestMapping("/spu")
public String getProductStatsGroupBySpu(
@RequestParam(value = "date", defaultValue = "0") Integer date,
@RequestParam(value = "limit", defaultValue = "10") int limit) {
if (date == 0) date = now();
List<ProductStats> statsList
= productService.getProductStatsGroupBySpu(date, limit);
//設定表頭
StringBuilder jsonBuilder =
new StringBuilder(" " +
"{\"status\":0,\"data\":{\"columns\":[" +
"{\"name\":\"商品名稱\",\"id\":\"spu_name\"}," +
"{\"name\":\"交易額\",\"id\":\"order_amount\"}," +
"{\"name\":\"訂單數\",\"id\":\"order_ct\"}]," +
"\"rows\":[");
//回圈拼接表體
for (int i = 0; i < statsList.size(); i++) {
ProductStats productStats = statsList.get(i);
if (i >= 1) {
jsonBuilder.append(",");
}
jsonBuilder.append("{\"spu_name\":\"" + productStats.getSpu_name() + "\"," +
"\"order_amount\":" + productStats.getOrder_amount() + "," +
"\"order_ct\":" + productStats.getOrder_ct() + "}");
}
jsonBuilder.append("]}}");
return jsonBuilder.toString();
}
//品類介面方法
@RequestMapping("/category3")
public String getProductStatsGroupByCategory3(
@RequestParam(value = "date", defaultValue = "0") Integer date,
@RequestParam(value = "limit", defaultValue = "4") int limit) {
if (date == 0) {
date = now();
}
List<ProductStats> statsList
= productService.getProductStatsGroupByCategory3(date, limit);
StringBuilder dataJson = new StringBuilder("{ \"status\": 0, \"data\": [");
int i = 0;
for (ProductStats productStats : statsList) {
if (i++ > 0) {
dataJson.append(",");
}
;
dataJson.append("{\"name\":\"")
.append(productStats.getCategory3_name()).append("\",");
dataJson.append("\"value\":")
.append(productStats.getOrder_amount()).append("}");
}
dataJson.append("]}");
return dataJson.toString();
}
//品牌介面方法
@RequestMapping("/trademark")
public String getProductStatsByTrademark(
@RequestParam(value = "date", defaultValue = "0") Integer date,
@RequestParam(value = "limit", defaultValue = "10") int limit) {
if (date == 0) {
date = now();
}
List<ProductStats> productStatsByTrademarkList
= productService.getProductStatsByTrademark(date, limit);
List<String> tradeMarkList = new ArrayList<>();
List<BigDecimal> amountList = new ArrayList<>();
for (ProductStats productStats : productStatsByTrademarkList) {
tradeMarkList.add(productStats.getTm_name());
amountList.add(productStats.getOrder_amount());
}
String json = "{\"status\":0,\"data\":{" + "\"categories\":" +
"[\"" + StringUtils.join(tradeMarkList, "\",\"") + "\"],\"series\":[" +
"{\"data\":[" + StringUtils.join(amountList, ",") + "]}]}}";
return json;
}
@RequestMapping("/province")
public String getProvinceStats(@RequestParam(value = "date", defaultValue = "0") Integer date) {
if (date == 0) {
date = now();
}
StringBuilder jsonBuilder = new StringBuilder("{\"status\":0,\"data\":{\"mapData\":[");
List<ProvinceStats> provinceStatsList = productService.getProvinceStats(date);
if (provinceStatsList.size() == 0) {
// jsonBuilder.append( "{\"name\":\"北京\",\"value\":0.00}");
}
for (int i = 0; i < provinceStatsList.size(); i++) {
if (i >= 1) {
jsonBuilder.append(",");
}
ProvinceStats provinceStats = provinceStatsList.get(i);
jsonBuilder.append("{\"name\":\"" + provinceStats.getProvince_name() + "\",\"value\":" + provinceStats.getOrder_amount() + " }");
}
jsonBuilder.append("]}}");
return jsonBuilder.toString();
}
@RequestMapping("/visitor")
public String getVisitorStatsByNewFlag(@RequestParam(value = "date", defaultValue = "0") Integer date) {
if (date == 0) date = now();
List<VisitorStats> visitorStatsByNewFlag = productService.getVisitorStatsByNewFlag(date);
VisitorStats newVisitorStats = new VisitorStats();
VisitorStats oldVisitorStats = new VisitorStats();
//回圈把資料賦給新訪客統計物件和老訪客統計物件
for (VisitorStats visitorStats : visitorStatsByNewFlag) {
if (visitorStats.getIs_new().equals("1")) {
newVisitorStats = visitorStats;
} else {
oldVisitorStats = visitorStats;
}
}
//把資料拼接入字串
String json = "{\"status\":0,\"data\":{\"combineNum\":1,\"columns\":" +
"[{\"name\":\"類別\",\"id\":\"type\"}," +
"{\"name\":\"新用戶\",\"id\":\"new\"}," +
"{\"name\":\"老用戶\",\"id\":\"old\"}]," +
"\"rows\":" +
"[{\"type\":\"用戶數(人)\"," +
"\"new\": " + newVisitorStats.getUv_ct() + "," +
"\"old\":" + oldVisitorStats.getUv_ct() + "}," +
"{\"type\":\"總訪問頁面(次)\"," +
"\"new\":" + newVisitorStats.getPv_ct() + "," +
"\"old\":" + oldVisitorStats.getPv_ct() + "}," +
"{\"type\":\"跳出率(%)\"," +
"\"new\":" + newVisitorStats.getUjRate() + "," +
"\"old\":" + oldVisitorStats.getUjRate() + "}," +
"{\"type\":\"平均在線時長(秒)\"," +
"\"new\":" + newVisitorStats.getDurPerSv() + "," +
"\"old\":" + oldVisitorStats.getDurPerSv() + "}," +
"{\"type\":\"平均訪問頁面數(人次)\"," +
"\"new\":" + newVisitorStats.getPvPerSv() + "," +
"\"old\":" + oldVisitorStats.getPvPerSv()
+ "}]}}";
return json;
}
@Autowired
KeywordService keywordService;
@RequestMapping("/keyword")
public String getKeywordStats(@RequestParam(value = "date",defaultValue = "0") Integer date,
@RequestParam(value = "limit",defaultValue = "20") int limit){
if(date==0){
date=now();
}
//查詢資料
List<KeywordStats> keywordStatsList
= keywordService.getKeywordStats(date, limit);
StringBuilder jsonSb=new StringBuilder( "{\"status\":0,\"msg\":\"\",\"data\":[" );
//回圈拼接字串
for (int i = 0; i < keywordStatsList.size(); i++) {
KeywordStats keywordStats = keywordStatsList.get(i);
if(i>=1){
jsonSb.append(",");
}
jsonSb.append( "{\"name\":\"" + keywordStats.getKeyword() + "\"," +
"\"value\":"+keywordStats.getCt()+"}");
}
jsonSb.append( "]}");
return jsonSb.toString();
}
}
bean
keywordStats
package com.example.gmallpublisher.bean;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* Desc: 關鍵詞統計物體類
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class KeywordStats {
private String keyword;
private Long ct;
private String source;
private String stt;
private String edt;
private Long ts;
}
ProductStats
package com.example.gmallpublisher.bean;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
/**
* Desc: 商品交易額統計物體類
*/
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class ProductStats {
String stt;
String edt;
Long sku_id;
String sku_name;
BigDecimal sku_price;
Long spu_id;
String spu_name;
Long tm_id ;
String tm_name;
Long category3_id ;
String category3_name ;
@Builder.Default
Long display_ct=0L;
@Builder.Default
Long click_ct=0L;
@Builder.Default
Long cart_ct=0L;
@Builder.Default
Long order_sku_num=0L;
@Builder.Default
BigDecimal order_amount=BigDecimal.ZERO;
@Builder.Default
Long order_ct=0L;
@Builder.Default
BigDecimal payment_amount=BigDecimal.ZERO;
@Builder.Default
Long refund_ct=0L;
@Builder.Default
BigDecimal refund_amount=BigDecimal.ZERO;
@Builder.Default
Long comment_ct=0L;
@Builder.Default
Long good_comment_ct=0L ;
Long ts;
}
ProvinceStats
package com.example.gmallpublisher.bean;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
/**
* Desc: 地區交易額統計物體類
*/
@AllArgsConstructor
@Data
@NoArgsConstructor
public class ProvinceStats {
private String stt;
private String edt;
private String province_id;
private String province_name;
private BigDecimal order_amount;
private String ts;
}
VisitorStats
package com.example.gmallpublisher.bean;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.math.BigDecimal;
import java.math.RoundingMode;
/**
* Desc: 訪客流量統計物體類
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class VisitorStats {
private String stt;
private String edt;
private String vc;
private String ch;
private String ar;
private String is_new;
private Long uv_ct = 0L;
private Long pv_ct = 0L;
private Long sv_ct = 0L;
private Long uj_ct = 0L;
private Long dur_sum = 0L;
private Long new_uv = 0L;
private Long ts;
private int hr;
//計算跳出率 = 跳出次數*100/訪問次數
public BigDecimal getUjRate() {
if (uv_ct != 0L) {
return BigDecimal.valueOf(uj_ct)
.multiply(BigDecimal.valueOf(100))
.divide(BigDecimal.valueOf(sv_ct), 2, RoundingMode.HALF_UP);
} else {
return BigDecimal.ZERO;
}
}
//計算每次訪問停留時間(秒) = 當日總停留時間(毫秒)/當日訪問次數/1000
public BigDecimal getDurPerSv() {
if (uv_ct != 0L) {
return BigDecimal.valueOf(dur_sum)
.divide(BigDecimal.valueOf(sv_ct), 0, RoundingMode.HALF_UP)
.divide(BigDecimal.valueOf(1000), 1, RoundingMode.HALF_UP);
} else {
return BigDecimal.ZERO;
}
}
//計算每次訪問停留頁面數 = 當日總訪問頁面數/當日訪問次數
public BigDecimal getPvPerSv() {
if (uv_ct != 0L) {
return BigDecimal.valueOf(pv_ct)
.divide(BigDecimal.valueOf(sv_ct), 2, RoundingMode.HALF_UP);
} else {
return BigDecimal.ZERO;
}
}
}
最終效果

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/296764.html
標籤:其他
上一篇:2021-08-29 兩數相加
下一篇:二叉樹展開為鏈表 C語言
