==[案例]Spark實(shí)時(shí)統(tǒng)計(jì)訂單量

Spark實(shí)時(shí)統(tǒng)計(jì)訂單量 - 簡書
http://www.reibang.com/p/3ec093a9d584

Paste_Image.png

注:組件不了解的同學(xué)可參考其他文章,本文主要講項(xiàng)目的實(shí)現(xiàn)
1、某些同學(xué)會(huì)問,直接在業(yè)務(wù)系統(tǒng)加入JS埋點(diǎn)通過發(fā)日志不更好嗎?
答:第一、JS埋點(diǎn)業(yè)務(wù)系統(tǒng)涉及產(chǎn)品改造,不可能因?yàn)橐粋€(gè)需求讓你去隨便改業(yè)務(wù)系統(tǒng)醇滥。第二、即使加入JS埋點(diǎn)也不可能獲得業(yè)務(wù)系統(tǒng)的全部數(shù)據(jù)超营。所以業(yè)務(wù)系統(tǒng)核心數(shù)據(jù)還得去業(yè)務(wù)系統(tǒng)庫獲取鸳玩。

2、還有人問加入Kafka太多余
答:第一糟描、加入Kafka為了使系統(tǒng)擴(kuò)展性更強(qiáng)怀喉,可方便對(duì)接各種開源產(chǎn)品。第二船响、通過Kafka消息組可使同一條消息被不同Consumer消費(fèi)躬拢,用戶離線和實(shí)時(shí)兩條線。


前言
本人GitHub地址:https://github.com/guofei1219QQ : 86608625咨詢項(xiàng)目相關(guān)問題的請(qǐng)直接說明問題见间,不要一直問在嗎聊闯?還在嗎?等問題米诉,博主QQ一直健在呢菱蔬,由于本人平時(shí)還要工作,問題不能及時(shí)回復(fù)請(qǐng)見諒J仿隆K┟凇!
背景
用戶下單數(shù)據(jù)會(huì)通過業(yè)務(wù)系統(tǒng)實(shí)時(shí)產(chǎn)生入庫到mysql庫惊橱,我們要統(tǒng)計(jì)通某個(gè)推廣渠道實(shí)時(shí)下單量蚪腐,以便線上運(yùn)營推廣人員查看不同渠道推廣效果進(jìn)而執(zhí)行不同推廣策略
系統(tǒng)架構(gòu)

架構(gòu)圖

注:組件不了解的同學(xué)可參考其他文章,本文主要講項(xiàng)目的實(shí)現(xiàn)1税朴、某些同學(xué)會(huì)問回季,直接在業(yè)務(wù)系統(tǒng)加入JS埋點(diǎn)通過發(fā)日志不更好嗎家制?答:第一、JS埋點(diǎn)業(yè)務(wù)系統(tǒng)涉及產(chǎn)品改造泡一,不可能因?yàn)橐粋€(gè)需求讓你去隨便改業(yè)務(wù)系統(tǒng)颤殴。第二、即使加入JS埋點(diǎn)也不可能獲得業(yè)務(wù)系統(tǒng)的全部數(shù)據(jù)鼻忠。所以業(yè)務(wù)系統(tǒng)核心數(shù)據(jù)還得去業(yè)務(wù)系統(tǒng)庫獲取涵但。
2、還有人問加入Kafka太多余答:第一粥烁、加入Kafka為了使系統(tǒng)擴(kuò)展性更強(qiáng)贤笆,可方便對(duì)接各種開源產(chǎn)品。第二讨阻、通過Kafka消息組可使同一條消息被不同Consumer消費(fèi),用戶離線和實(shí)時(shí)兩條線篡殷。
解析Mysql binlog日志
主要邏輯
1.創(chuàng)建Canal連接2.解析Mysql binlog獲得insert語句
public static void main(String args[]) { //第一個(gè)參數(shù)為Canal server服務(wù)IP地址如果使用windows開發(fā)連接linux Canal服務(wù)需要制定IP eg: new InetSocketAddress("192.168.61.132", 11111) //第二個(gè)參數(shù)為Canal server服務(wù)端口號(hào) Canal server IP和端口號(hào)在 /conf/canal.properties中配置 //第三個(gè)參數(shù)為Canal instance名稱 /conf下目錄名稱 //第四第五個(gè)參數(shù)為mysql用戶名和密碼钝吮,如果在 /conf/example/instance.properties中已經(jīng)配置 這里不用謝 CanalConnector connector = CanalConnectors.newSingleConnector(new InetSocketAddress("192.168.61.132", 11111), "example", "", ""); int batchSize = 1000; int emptyCount = 0; try { connector.connect(); connector.subscribe(".\.."); connector.rollback(); int totalEmtryCount = 120; while (emptyCount < totalEmtryCount) { Message message = connector.getWithoutAck(batchSize); // 獲取指定數(shù)量的數(shù)據(jù) long batchId = message.getId(); int size = message.getEntries().size(); if (batchId == -1 || size == 0) { emptyCount++; System.out.println("empty count : " + emptyCount); try { Thread.sleep(1000); } catch (InterruptedException e) { } } else { emptyCount = 0; // System.out.printf("message[batchId=%s,size=%s] \n", batchId, size); printEntry(message.getEntries()); } connector.ack(batchId); // 提交確認(rèn) // connector.rollback(batchId); // 處理失敗, 回滾數(shù)據(jù) } System.out.println("empty too many times, exit"); } finally { connector.disconnect(); }}

組裝數(shù)據(jù)發(fā)送至Kafka
private static void printColumn(List<Column> columns) { for (Column column : columns) { System.out.println(column.getName() + " : " + column.getValue()); KafkaProducer.sendMsg("canal", UUID.randomUUID().toString() ,column.getName() + " : " + column.getValue()); }}

Streaming分渠道匯總數(shù)據(jù)
以DStream中的數(shù)據(jù)進(jìn)行按key做reduce操作,然后對(duì)各個(gè)批次的數(shù)據(jù)進(jìn)行累加在有新的數(shù)據(jù)信息進(jìn)入或更新時(shí)板辽,可以讓用戶保持想要的任何狀奇瘦。使用這個(gè)功能需要完成兩步:1) 定義狀態(tài):可以是任意數(shù)據(jù)類型2) 定義狀態(tài)更新函數(shù):用一個(gè)函數(shù)指定如何使用先前的狀態(tài),從輸入流中的新值更新狀態(tài)劲弦。對(duì)于有狀態(tài)操作耳标,要不斷的把當(dāng)前和歷史的時(shí)間切片的RDD累加計(jì)算,隨著時(shí)間的流失邑跪,計(jì)算的數(shù)據(jù)規(guī)模會(huì)變得越來越大次坡。
val orders = resut_lines.updateStateByKey(updateRunningSum _)def updateRunningSum(values: Seq[Long], state: Option[Long]) = {/* state:存放的歷史數(shù)據(jù) values:當(dāng)前批次匯總值 */Some(state.getOrElse(0L)+values.sum)}

統(tǒng)計(jì)結(jié)果寫入Mysql
實(shí)時(shí)匯總某渠道下單量需要根據(jù)渠道為主鍵更新或插入新數(shù)據(jù)1.當(dāng)某個(gè)渠道第一單時(shí),庫中沒有以此渠道為主鍵的數(shù)據(jù)画畅,需要insert into 訂單統(tǒng)計(jì)表2.當(dāng)某渠道在庫中已有該渠道下單量砸琅,需要更新此渠道下單量值 update 訂單統(tǒng)計(jì)表所以我們使用:

有該渠道就更新,沒有就插入REPLACE INTO order_statistic(chanel, orders) VALUES(?, ?)

orders.foreachRDD(rdd =>{ rdd.foreachPartition(rdd_partition =>{ rdd_partition.foreach(data=>{ if(!data.toString.isEmpty) { System.out.println("訂單量"+" : "+data._2) DataUtil.toMySQL(data._1.toString,data._2.toInt) } }) })})def toMySQL(name: String,orders:Int) = { var conn: Connection = null var ps: PreparedStatement = null val sql = "REPLACE INTO order_statistic(chanel, orders) VALUES(?, ?)" try { Class.forName("com.mysql.jdbc.Driver"); conn = DriverManager.getConnection("jdbc:mysql://192.168.20.126:3306/test", "root", "root") ps = conn.prepareStatement(sql) ps.setString(1, name) ps.setInt(2, orders) ps.executeUpdate() } catch { case e: Exception => e.printStackTrace() } finally { if (ps != null) { ps.close() } if (conn != null) { conn.close() } }}

FAQ
1.canal依賴Canal protobuf版本為2.4.1轴踱,而spark依賴的2.5版本
<dependency> <groupId>com.google.protobuf</groupId> <artifactId>protobuf-java</artifactId> <version>2.4.1</version></dependency>

參考文章
1.Canal wiki:https://github.com/alibaba/canal/wiki2.streaming關(guān)于轉(zhuǎn)化操作http://spark.apache.org/docs/1.6.0/streaming-programming-guide.html#transformations-on-dstreams3.mysql的replace intohttp://blog.sina.com.cn/s/blog_5f53615f01016wy3.html
文/MichaelFly(簡書作者)原文鏈接:http://www.reibang.com/p/3ec093a9d584著作權(quán)歸作者所有症脂,轉(zhuǎn)載請(qǐng)聯(lián)系作者獲得授權(quán),并標(biāo)注“簡書作者”淫僻。

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末诱篷,一起剝皮案震驚了整個(gè)濱河市,隨后出現(xiàn)的幾起案子雳灵,更是在濱河造成了極大的恐慌棕所,老刑警劉巖,帶你破解...
    沈念sama閱讀 218,858評(píng)論 6 508
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件细办,死亡現(xiàn)場離奇詭異橙凳,居然都是意外死亡蕾殴,警方通過查閱死者的電腦和手機(jī),發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 93,372評(píng)論 3 395
  • 文/潘曉璐 我一進(jìn)店門岛啸,熙熙樓的掌柜王于貴愁眉苦臉地迎上來钓觉,“玉大人,你說我怎么就攤上這事坚踩〉丛郑” “怎么了?”我有些...
    開封第一講書人閱讀 165,282評(píng)論 0 356
  • 文/不壞的土叔 我叫張陵瞬铸,是天一觀的道長批幌。 經(jīng)常有香客問我,道長嗓节,這世上最難降的妖魔是什么荧缘? 我笑而不...
    開封第一講書人閱讀 58,842評(píng)論 1 295
  • 正文 為了忘掉前任,我火速辦了婚禮拦宣,結(jié)果婚禮上截粗,老公的妹妹穿的比我還像新娘。我一直安慰自己鸵隧,他們只是感情好绸罗,可當(dāng)我...
    茶點(diǎn)故事閱讀 67,857評(píng)論 6 392
  • 文/花漫 我一把揭開白布。 她就那樣靜靜地躺著豆瘫,像睡著了一般珊蟀。 火紅的嫁衣襯著肌膚如雪。 梳的紋絲不亂的頭發(fā)上外驱,一...
    開封第一講書人閱讀 51,679評(píng)論 1 305
  • 那天育灸,我揣著相機(jī)與錄音,去河邊找鬼略步。 笑死描扯,一個(gè)胖子當(dāng)著我的面吹牛,可吹牛的內(nèi)容都是我干的趟薄。 我是一名探鬼主播绽诚,決...
    沈念sama閱讀 40,406評(píng)論 3 418
  • 文/蒼蘭香墨 我猛地睜開眼,長吁一口氣:“原來是場噩夢(mèng)啊……” “哼杭煎!你這毒婦竟也來了恩够?” 一聲冷哼從身側(cè)響起,我...
    開封第一講書人閱讀 39,311評(píng)論 0 276
  • 序言:老撾萬榮一對(duì)情侶失蹤羡铲,失蹤者是張志新(化名)和其女友劉穎蜂桶,沒想到半個(gè)月后,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體也切,經(jīng)...
    沈念sama閱讀 45,767評(píng)論 1 315
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡扑媚,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 37,945評(píng)論 3 336
  • 正文 我和宋清朗相戀三年腰湾,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片疆股。...
    茶點(diǎn)故事閱讀 40,090評(píng)論 1 350
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡费坊,死狀恐怖,靈堂內(nèi)的尸體忽然破棺而出旬痹,到底是詐尸還是另有隱情附井,我是刑警寧澤,帶...
    沈念sama閱讀 35,785評(píng)論 5 346
  • 正文 年R本政府宣布两残,位于F島的核電站永毅,受9級(jí)特大地震影響,放射性物質(zhì)發(fā)生泄漏人弓。R本人自食惡果不足惜沼死,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 41,420評(píng)論 3 331
  • 文/蒙蒙 一、第九天 我趴在偏房一處隱蔽的房頂上張望票从。 院中可真熱鬧漫雕,春花似錦、人聲如沸峰鄙。這莊子的主人今日做“春日...
    開封第一講書人閱讀 31,988評(píng)論 0 22
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽吟榴。三九已至,卻和暖如春囊扳,著一層夾襖步出監(jiān)牢的瞬間吩翻,已是汗流浹背。 一陣腳步聲響...
    開封第一講書人閱讀 33,101評(píng)論 1 271
  • 我被黑心中介騙來泰國打工锥咸, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留狭瞎,地道東北人。 一個(gè)月前我還...
    沈念sama閱讀 48,298評(píng)論 3 372
  • 正文 我出身青樓搏予,卻偏偏與公主長得像熊锭,于是被迫代替她去往敵國和親。 傳聞我的和親對(duì)象是個(gè)殘疾皇子雪侥,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 45,033評(píng)論 2 355

推薦閱讀更多精彩內(nèi)容