一、Clickhouse 介紹
1.1 Clickhouse 介紹
ClickHouse 是一個用于聯機分析(OLAP)的列式資料庫管理系統(DBMS),最初 是一款名為 Yandex.Metrica 的產品,主要用于 WEB 流量分析,ClickHouse 的全稱是 Click Stream,Data WareHouse,簡稱 ClickHouse,
1.2 Clickhouse 的分布式架構

1.3 Clickhouse 的特性
(1)真正的面向列的DBMS
在一個真正的面向列的DBMS中,沒有任何“垃圾”存盤在值中,例如,必須支持定長數值,以避免在數值旁邊存盤長度“數字”,例如,十億個UInt8型別的值實際_上應該消耗大約1 GB的未壓縮磁盤空間,否則這將強烈影響CPU的使用,由于解壓縮的速度(CPU 使用率)主要取決于未壓縮的資料量,所以即使在未壓縮的情況下,緊湊地存盤資料(沒有任何“垃圾”)也是非常重要的,另外,ClickHouse 是- -個DBMS,而不是一一個單- -的資料庫,ClickHouse 允許在運行時創建表和資料庫,加載資料和運行查詢,而無需重新配置和重新啟動服務器,
(2)資料壓縮
資料壓縮可以提高性能,ClickHouse除了在磁盤空間和CPU消耗之間進行不同權衡的高效通用壓縮編解碼器之外,ClickHouse還提供針對特定型別資料的專用編解碼器,這使得ClickHouse能夠與更小的資料庫(如時間序列資料庫)競爭并 超越它們,
(3)磁盤存盤的資料
ClickHouse被設計用于作業在傳統磁盤上的系統,它提供每GB更低的存盤成本,但如果可以使用SSD和記憶體,它也會合理的利用這些資源,
(4)多核并行處理
多核多節點并行化大型查詢,
(5)在多個服務器上分布式處理
在ClickHouse中,資料可以駐留在不同的分片上,每個分片可以是用于容錯的一組副本,查詢在所有分片上并行處理,
(6)SQL支持
(7)向量化引擎
資料不僅按列存盤,而且由矢量-列的部分進行處理,這使我們能夠提高高CPU使用性能,
(8)實時資料更新
ClickHouse支持主鍵表,為了快速執行對主鍵范圍的查詢,資料使用合并樹(MergeTree)以增量的方式有序的存盤,由于這個原因,資料可以不斷地添加到表中,添加資料時無鎖處理,
(9)索引
按照主鍵對資料進行排序,這將幫助ClickHouse在幾十毫秒以內完成對資料特定值或范圍的查找,
(10)支持在線查詢
在沒有對資料做任何預處理的情況下以極低的延遲處理查詢并將結果加載到用戶的頁面中,
(11)支持近似計算
- 系統包含用于近似計算各種值,中位數和分位數的集合函式,
- 支持基于部分(樣本)資料運行查詢并獲得近似結果,在這種情況下,從磁盤檢索比例較少的資料,
- 支持為有限數量的隨機密鑰(而不是所有密鑰)運行聚合,在資料中密鑰分發的特定條件下,這提供了相對準確的結果,同時使用較少的資源
(12)資料復制和對資料完整性的支持
使用異步多主復制,寫入任何可用的副本后,資料將分發到所有剩余的副本,系統在不同的副本上保持相同的資料,資料在失敗后自動恢復,在一些少數的復雜情況下需要手動恢復,
1.4 ClickHouse的缺點
(1)不支持事務,
(2)不擅長Update/Delete操作,僅能用于批量洗掉或修改資料,
(3)不擅長根據主鍵按行粒度進行查詢( 雖然支持),故不應該把ClickHouse當作Key-Value資料庫使用,
二、ClickHouse 的資料型別

2.1 整型
固定長度的整型,包括有符號整型或無符號整型,
整型范圍(-2^n-1 ~2^n-1-1) :
Int8 -[-128: 127]
Int16- [-32768 : 32767]
Int32 - [-2147483648 : 2147483647]
Int64 - [-9223372036854775808 : 9223372036854775807]
無符號整型范圍(0~2^n-1): 沒有負數,
UInt8- [0 : 255]
UInt16- [0 : 65535]
UInt32- [0 : 4294967295]
UInt64- [0 : 18446744073709551615]
2.2 浮點型
Float32 - float
Float64 - double
建議盡可能以整數形式存盤資料,例如,將固定精度的數字轉換為整數值,如時間用毫秒為單位表示,因為浮點型進行計算時可能引起四舍五入的誤差,

inf 正無窮

-inf 負無窮

nan 非數字

2.3 布爾型別(沒有)
clickhouse沒有布爾型別,可以使用UInt8 型別,利用列舉將取值限制為0或1,
2.4 字串
String
字串可以任意長度的,它可以包含任意的位元組集,包含空位元組,
FixedString(N)
固定長度N的字串,N必須是嚴格的正自然數,當服務端讀取長度小于N的字串時候,通過在字串末尾添加空位元組來達到N位元組長度,當服務端讀取長度大于N的字串時候,將回傳錯誤訊息,
與String 相比,極少會使用FixedString, 因為使用起來不是很方便,
2.5 列舉
包括Enum8和Enum16型別,Enum只支持'string'= int的對應型別,
Enum8用'String'= Int8對描述,
Enum16用'String'= Int16對描述,

可以通過 cast 轉換函式,列印出來對應的值,

2.6 陣列
Array(T):由T型別元素組成的陣列,
T可以是任意型別,包含陣列型別,但不推薦使用多維陣列,ClickHouse 對多維陣列的支持有限,例如,不能在MergeTree表中存盤多維陣列,
創建陣列方式1,使用array 函式:

創建陣列方式 2:使用方括號:

2.7 元組
創建元組方式 1,使用 tuple 函式:

創建元組方式 2,使用()即可:

2.8 日期
目前ClickHouse 有三種時間型別
Date接受年月-日的字串比如2019-12-16'
Datetime接受年月-日時:分:秒的字串 比如2019-12-16 20:50:10'
Datetime64接受年一月一日時:分:秒,亞秒的字串比如' 2019-12-16 20:50:10.66'
所有的時間日期函式都可以在第二個可選引數中接受時區引數,示例: Asia/ Yekaterinburg,在這種情況下,它們使用指定的時區而不是本地(默認)時區,
SELECT
toDateTime('2021-01-01 23:00:00') AS time,
toDate(time) AS date_local,
toDate(time, 'Asia/Yekaterinburg') AS date_yekat,
toString(time, 'US/Samoa') AS time_samoa;

常用的日期處理函式:
now() # 2020-04-01 17:25:40 取當前時間
toYear() # 2020 取日期中的年份
toMonth() # 4 取日期中的月份
today() # 2020-04-01 今天的日期
yesterday() # 2020-03-31 昨天的額日期
toDayOfYear() # 92 取一年中的第幾天
toDayOfWeek() # 3 取一周中的第幾天
toHour() # 17 取小時
toMinute() # 25 取分鐘
toSecond() # 40 取秒
toStartOfYear() # 2020-01-01 取一年中的第一天
toStartOfMonth() # 2020-04-01 取當月的第一天
formatDateTime(now(),'%Y-%m-%d') # 2020*04-01 指定時間格式
toYYYYMM() # 202004
toYYYYMMDD() # 20200401
toYYYYMMDDhhmmss() # 20200401172540
dateDiff()
......
三、表引擎
3.1 表引擎的介紹
表引擎是ClickHouse 的一大特色,如果對MySQL熟悉的話,或許你應該聽說過InnoDB,可以說,表引 擎決定了如何存盤表的資料,
包括:
資料的存盤方式和位置,寫到哪里以及從哪里讀取資料,
支持哪些查詢以及如何支持,
并發資料訪問,
索引的使用(如果存在),
是否可以執行多執行緒請求,
資料復制引數,
表引擎的使用方式就是必須顯式在創建表時定義該表使用的引擎,以及引擎.使用的相關引數,
特別注意:引擎的名稱大小寫敏感,
3.2 表引擎的分類

3.3 Log系串列引擎
3.3.1 Log系串列引擎的介紹
Log系串列引擎功能相對簡單,主要用于快速寫入小表(1百萬行左右的表), 然后全部讀出的場景,即一次寫入多次查詢,
3.3.2 Log系串列引擎的特點
資料存盤在磁盤上,
當寫資料時,將資料追加到檔案的末尾,
不支持并發讀寫,當向表中寫入資料時,針對這張表的查詢會被阻塞,直至寫入動作結束,
不支持索引,
不支持原子寫:如果某些操作(例外的服務器關閉)中斷了寫操作,則可能會獲得帶有損壞資料的表,
不支持ALTER操作(這些操作會修改表設定或資料,比如delete、update 等等),
3.3.3 Log系串列引擎的區別
TinyLog是Log系列引擎中功能簡單、性能較低的引擎,它的存盤結構由資料檔案和元資料兩部分組成,其中,資料檔案是按列獨立存盤的,也就是說每一個列欄位都對應一一個檔案, 除此之外,TinyLog 不支持并發資料讀取,
StripLog支持并發讀取資料檔案,當讀取資料時,ClickHouse 會使用多執行緒進行讀取,每個執行緒處理一個單獨的資料塊,另外,StripLog 將所有列資料存盤Log支持并發讀取資料檔案,當讀取資料時,ClickHouse會使用多執行緒進行讀取,每個執行緒處理-~個單獨的資料塊,Log引擎會將每個列資料單獨存盤在一個獨立檔案中,
3.3.4 TinyLog
該引擎適用于一次寫入,多次讀取的場景,對于處理小批資料的中間表可以使用該引擎,值得注意的是,使用大量的小表存盤資料,性能會很低,
CREATE TABLE emp_tinylog (
emp_id UInt16 COMMENT'員工id',
name String COMMENT'員工姓名',
work_place String COMMENT'作業地點',
age UInt8 COMMENT'員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=TinyLog();
INSERT INTO emp_tinylog
VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
INSERT INTO emp_tinylog
VALUES (3,'bob','北京',33,'財務部',50000),(4,'tony','杭州',28,'銷售事部',50000);
可以在/var/lib/clickhouse/data查看CK中的資料,TinyLog引擎表每一列都對應的檔案,在sizes.json檔案內使用JSON格式記錄了每個. bin 檔案內對應的資料大小的資訊,

當我們執行ALTER操作時會報錯,說明該表引擎不支持ALTER操作,
ALTER TABLE emp_tinylog DELETE WHERE emp_id= 5;
ALTER TABLE emp_inylog UPDATE age = 30 WHERE emp_id= 4;

3.3.5 StripLog
相比TinyLog而言,StripeLog 擁有更高的查詢性能( 擁有.mrk標記檔案,支持并行查詢),同時其使用了更少的檔案描述符(所有資料使用同一個檔案保存),
CREATE TABLE emp_stripelog (
emp_id UInt16 COMMENT '員工id',
name String COMMENT'員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=StripeLog;
--插入資料
INSERT INTO emp_stripelog
VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
INSERT INTO emp_stripelog
VALUES (3,'bob','北京',33,'財務部',0000),(,'tony','杭州',28,'銷售事部',50000);
--查詢資料
select * from emp_stripelog;

進入默認資料存盤目錄,查看底層資料存盤形式,

可以看出StripeLog表引擎對應的存盤結構包括三個檔案:
data.bin:資料檔案,所有的列欄位使用同一個檔案保存,它們的資料都會被寫入data.bin.
index.mrk:資料標記,保存了資料在data.bin檔案中的位置資訊(每個插入資料塊對應列的offset),利用資料標記能夠使用多個執行緒,以并行的方式讀取data.bin內的壓縮資料塊,從而提升資料查詢的性能,
sizes.json:元資料檔案,記錄了data.bin和index.mrk大小的資訊,
注意: StripeLog 引擎將所有資料都存盤在了一個檔案中,對于每次的INSERT操作,ClickHouse 會將資料塊追加到表檔案的末尾,StripeLog 引擎同樣不支持ALTER UPDATE和ALTER DELETE操作,
3.3.6 Log
Log引擎表適用于臨時資料,一次性寫入、測驗場景,Log引擎結合了TinyLog ;表引擎和Stripelog表引擎的長處,是Log系列引擎中性能最高的表引擎,
CREATE TABLE emp_log(
emp_id UInt16 COMMENT '員工id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=Log;
-- 插入數就
INSERT INTO emp_log VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',1000);
INSERT INTO emp_log VALUES (3,'bob','北京',33,'財務部',50000),(4,'tony','杭州',28,'銷售事部',50000);
Log 引擎的存盤結構:

_marks.mrk:資料標記,統一保存了資料在各個.bin檔案中的位置資訊,利用資料標記能夠使用多個執行緒,以并行的方式讀取,.bin 內的壓縮資料塊,從而提升資料查詢的性能,Log 表引擎會將每一列都存在-一個檔案中,對于每一次的INSERT操作,都會對應一個資料塊,
3.4 MergeTree系串列引擎
在所有的表引擎中,最為核心的當屬MergeTree系串列引擎,這些表引擎擁有最為強大的性能和最廣泛的使用場合,對于非MergeTree系列的其他引擎而言,主要用于特殊用途,場景相對有限,而MergeTree系串列引擎是官方主推的存盤引擎,支持幾乎所有ClickHouse核心功能,
3.4.1 MergeTree
MergeTree在寫入一批資料時,資料總會以資料片段的形式寫入磁盤,且資料片段不可修改,為了避免片段過多,ClickHouse 會通過后臺執行緒,定期合并這些資料片段,屬于相同磁區的資料片段會被合成-一個新的片段,這種資料片段往復合并的特點,也正是合并樹名稱的由來,
特點:
需要指定主鍵,資料按照主鍵排序,
可以使用磁區,可以通過PARTITION KEY陳述句指定磁區欄位,開發中一般按月進行磁區,
支持資料副本,保證安全性,
支持資料采樣,
格式:

ENGINE - 引擎名和引數,MergeTree 無引數,
ORDERBY - 排序鍵,可以是一-組列的元組或任意的運算式,例如: ORDER BY
(CounterlD, EventDate),
如果沒有使用PRIMARYKEY顯式指定的主鍵,ClickHouse會使用排序鍵作為主鍵,
PARTITIONBY - 磁區鍵,可選項,
要按月磁區,可以使用運算式toYYYMM(date_ column) ,這里的date_ column 是一個Date 型別的列,磁區名的格式會是"YYYMM",
PRIMARY KEY - 如果要選擇與排序鍵不同的主鍵,在這里指定,可選項,
默認情況下主鍵跟排序鍵(由ORDERBY子句指定)相同,
因此,大部分情況下不需要再專門指定-一個PRIMARY KEY子句,
SAMPLE BY - 用于抽樣的運算式,可選項,
如果要用抽樣運算式,主鍵中必須包含這個運算式,例如: .
SAMPLE BY intHash32(UserlD) ORDER BY (CounterlD, EventDate,intHash32(UserlD)),
TTL - 資料的存活時間,在MergeTree中,可以為某個列欄位或整張表設定TTL,當時間到達時,如果是列欄位級別的TTL,則會洗掉這一-列的資料;如果是表級別的TTL,則會洗掉整張表的資料,可選項,
SETTINGS - 控制MergeTree 行為的額外引數,可選項:
index_ granularity — 索引粒度,索引中相鄰的 【標記】間的資料行數,默認值8192
use_ minimalistic _part_ header_in_zookeeper 一 ZooKeeper 中資料片段存盤方式,如果use_minimalistic_part_header_in_zookeeper=1 ,ZooKeeper 會存盤更少的資料,
min_merge_bytes_to_use_direct_io — 使用直接I/0 來操作磁盤的合并操作時要求的最小資料量,合并資料片段時,ClickHouse會計算要被合并的所有資料的總存盤空間,如果大小超過了min_merge_bytes_to_use_direct_io設定的位元組數,則ClickHouse將使用直接I/O 介面(O_DIRECT選項)對磁盤讀寫,如果設定min_ merge_ .bytes_ to_ use_ _direct_ io=0,則會禁用直接l/0, 默認值: 10* 1024* 1024* 1024位元組,(大于10G直接IO,小于10G緩沖IO),
CREATE TABLE emp_mergetree (
emp_id UInt16 COMMENT'員工id',
name String COMMENT'員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT'部門',
salary Decimal32(2) COMMENT'工資'
)ENGINE=MergeTree()
ORDER BY emp_id
PARTITION BY work_place;
--插入資料
INSERT INTO emp_mergetree VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
INSERT INTO emp_mergetree VALUES (3,'bob','北京',33,'財務部',50000),(4,'tony','杭州',28,'銷售事部',50000);
-- 查詢資料
select * from emp_mergetree;

可以在對應的/var/lib/clickhouse/data/test/emp_mergetree 表目錄下看到插入的三條資料被放入了不同子目錄中,

任何一個批次的資料寫入都會產生一個臨時磁區,不會納入任何一個已有的磁區,寫入后的某個時刻(大概10-15分鐘后),ClickHouse 會自動執行合并操作(等不及也可以手動通過optimize執行),把臨時磁區的資料,合并到已有磁區中,
--例如繼續插入資料:
INSERT INTO emp_mergetree VALUES (5,'robin','北京',35,'財務部',50000),(6,'lilei','北京',38,'銷售事部',50000);
發現資料沒有合并

手動觸發合并
optimize table emp_mergetree partition '北京';
再次查詢就能看到資料合并

插入相同主鍵,相同磁區資料
INSERT INTO emp_mergetree VALUES (1,'sam','上海',35,'財務部',50000);
觸發合并
optimize table emp_mergetree partition '上海';
再次查詢: (發現資料沒有去重,說明主鍵沒有去重功能,沒有唯一約束性)

在 MergeTree 中主鍵并不用于去重,而是用于索引,加快查詢速度,
3.4.2 ReplacingMergeTree
MergeTree表引擎無法對相同主鍵的資料進行去重,ClickHouse 提供了ReplacingMergeTree引擎,可以針對相同主鍵的資料進行去重,它能夠在合并磁區時洗掉重復的資料,值得注意的是,ReplacingMergeTree只是在一定程度上解決了資料重復問題,但是并不能完全保障資料不重復,
格式:

ver 一 版本列,型別為UInt*, Date或DateTime, 可選引數,
在資料合并的時候,ReplacingMergeTree 從所有具有相同排序鍵的行中選擇一行留下:
如果ver列未指定,保留最后一條,
如果ver 列已指定,保留ver值最大的版本,
示例:
CREATE TABLE emp_replacingmergetree (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=ReplacingMergeTree()
ORDER BY (emp_id,name)
PRIMARY KEY emp_id
PARTITION BY work_place;
--插入資料
INSERT INTO emp_replacingmergetree VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
INSERT INTO emp_replacingmergetree VALUES (3,'bob',' 北 京 ',33,' 財 務 部 ',50000),(4,'tony',' 杭 州 ',28,' 銷 售 事 部 ',50000);

插入相同磁區相同主鍵、相同磁區相同 ORDER BY 鍵的資料:
INSERT INTO emp_replacingmergetree VALUES (1,'susan','上海',26,'技術部',20000),(1,'tom','上海',26,'技術部',30000);
--手動觸發合并操作
optimize table emp_replacingmergetree final;

插入不同磁區相同 ORDER BY 鍵的資料:
INSERT INTO emp_replacingmergetree VALUES (1,'tom','北京',26,'技術部',40000);
--手動觸發合并操作:
optimize table emp_replacingmergetree final;

ReplacingMergeTree 是支持去重的,并且是按照 ORDERBY 排序鍵為基準進行 去重的,而不是主鍵,
- 通過測驗得到結論:
-
實際上是使用 order by 欄位作為唯一鍵
-
去重不能跨磁區
-
只有合并磁區才會進行去重
-
認定重復的資料保留版本欄位值最大的
-
如果版本欄位相同則按插入順序保留最后一筆
-
3.4.3 SummingMergeTree
該引擎繼承自MergeTree, 區別在于,當合并SummingMergeTree 表的資料片段時,ClickHouse會把所有具有相同主鍵的行合并為一-行,該行包含了被合并的行中具有數值資料型別的列的匯總值,即如果存在重復的資料,會對對這些重復的資料進行合并成- - 條資料,類似于group by的效果,
如果用戶只需要查詢資料的匯總結果,不關心明細資料,并且資料的匯總條:件是預先明確的,即GROUP BY的分組欄位是確定的,可以使用該表引擎,

columns-包含了將要被匯.總的列的列名的元組,叫選引數,
所選的列必須是數值型別,并且不可位于主鍵中,
如果沒有指定columns',ClickHouse會把所有不在主鍵中的數值型別的列都進行匯總,
示例:
CREATE TABLE emp_summingmergetree (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=SummingMergeTree(salary)
PARTITION BY work_place
ORDER BY (emp_id,name) -- 注意排序 key 是兩個欄位
PRIMARY KEY emp_id -- 主鍵是一個欄位
;
-- 插入資料
INSERT INTO emp_summingmergetree VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
INSERT INTO emp_summingmergetree VALUES (3,'bob',' 北 京 ',33,' 財 務 部 ',50000),(4,'tony',' 杭 州 ',28,' 銷 售 事 部 ',50000);

當我們再次插入具有相同 emp_id,name 的資料時
INSERT INTO emp_summingmergetree VALUES (1,'tom','上海',25,'資訊部',10000),(1,'tom','北京',26,'人事部',10000);
-- 執行合并操作
optimize table emp_summingmergetree final;

通過測驗得到結論:
SummingMergeTree 用 ORBER BY 排序鍵作為聚合資料的條件 Key,即如果排 序 key 是相同的,則會合并成一條資料,并對指定的合并欄位進行聚合,
以資料磁區為單位來聚合資料,當磁區合并時,同-一-資料磁區內聚合Key相同的資料會被合并匯總,而不同磁區之間的資料則不會被匯總,
如果沒有指定聚合欄位,則會按照非主鍵的數值型別欄位進行聚合,.
如果兩行資料除了排序欄位相同,其他的非聚合欄位不相同,那么在聚合發生時,會保留最初的那條資料,新插入的資料對應的那個欄位值會被舍棄,
3.4.4 Aggregatingmergetree
該表引擎繼承自MergeTree,可以使用AggregatingMergeTree 表來做增量資料統計聚合,如果要按一組規則來合 并減少行數,則使用AggregatingMergeTree是合適的,AggregatingMergeTree 是通過預先定義的聚合函式計算資料并通過二進制的格式存入表內,與SummingMergeTree的區別在于: SummingMergeTree對非主鍵列進行sum聚合,而AggregatingMergeTree則可以指定各種聚合函式,
示例:
CREATE TABLE emp_aggregatingmergeTree (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary AggregateFunction(sum,Decimal32(2)) COMMENT '工資' )ENGINE=AggregatingMergeTree()
PARTITION BY work_place
ORDER BY (emp_id,name) -- 注意排序 key 是兩個欄位
PRIMARY KEY emp_id -- 主鍵是一個欄位
;
對于AggregateFunction型別的列欄位,在進行資料的寫入和查詢時與其他的表引擎有很大區別,在寫入資料時,需要呼叫** -State函式;而在查詢資料時, 則需要呼叫相應的-Merge函式,對于上面的建表陳述句而言,需要使用sumState**函式進行資料插入,
-- 需要使用 INSERT…SELECT 陳述句進行資料插入
INSERT INTO TABLE emp_aggregatingmergeTree
SELECT 1,'tom','上海',25,'資訊部',sumState(toDecimal32(10000,2));
INSERT INTO TABLE emp_aggregatingmergeTree
SELECT 1,'tom','上海',25,'資訊部',sumState(toDecimal32(20000,2));
-- 查詢
SELECT emp_id, name , sumMerge(salary) FROM emp_aggregatingmergeTree GROUP BY emp_id,name;

上面演示的用法非常的麻煩,其實更多的情況下,我們可以結合物化視圖一 起使用,將它作為物化視圖的表引擎,而這里的物化視圖是作為其他資料表上層 的一種查詢視圖,
AggregatingMergeTree 通常作為物化視圖的表引擎,與普通 MergeTree 搭配使用,
-- 創建一個 MereTree 引擎的明細表
-- 用于存盤全量的明細資料
-- 對外提供實時查詢
CREATE TABLE emp_mergetree_base (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
)ENGINE=MergeTree()
ORDER BY (emp_id,name)
PARTITION BY work_place;
-- 創建一張物化視圖
-- 使用 AggregatingMergeTree 表引擎
CREATE MATERIALIZED VIEW view_emp_agg
ENGINE = AggregatingMergeTree()
PARTITION BY emp_id
ORDER BY (emp_id,name) AS SELECT emp_id, name, sumState(salary) AS salary
FROM emp_mergetree_base
GROUP BY emp_id,name;
-- 向基礎明細表 emp_mergetree_base 插入資料
INSERT INTO emp_mergetree_base VALUES (1,'tom','上海',25,'技術部',20000),(1,'tom','上海',26,'人事部',10000);
-- 查詢物化視圖
SELECT emp_id, name , sumMerge(salary)
FROM view_emp_agg
GROUP BY emp_id,name;

3.4.5 CollapsingMergeTree
CollapsingMergeTree就是一種通過以增代刪的思路,支持行級資料修改和洗掉的表引擎,它通過定義一.個sign標記位欄位,記錄資料行的狀態,如果sign標記為1,則表示這是一行有效的資料;如果sign標記為-1,則表示這行資料需要被洗掉,當CollapsingMergeTree磁區合并時,同一資料磁區內,sign 標記為1和-1的一組資料會被抵消洗掉,
每次需要新增資料時,寫入一行sign標記為1的資料;需要洗掉資料時,則寫入一行sign標記為-1的資料,

示例:
CREATE TABLE emp_collapsingmergetree (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資',
sign Int8
)ENGINE=CollapsingMergeTree(sign)
ORDER BY (emp_id,name)
PARTITION BY work_place;
CollapsingMergeTree同樣是以ORDER BY 排序鍵作為判斷資料唯一性的依據,
-- 插入新增資料,sign=1 表示正常資料
INSERT INTO emp_collapsingmergetree VALUES (1,'tom','上海',25,'技術部',20000,1);
-- 首先插入一條與原來相同的資料(ORDER BY 欄位一致),并將 sign 置為-1
INSERT INTO emp_collapsingmergetree VALUES (1,'tom','上海',25,'技術部',20000,-1);
-- 再插入更新之后的資料
INSERT INTO emp_collapsingmergetree VALUES (1,'tom','上海',25,'技術部',30000,1);
-- 查詢資料
select * from emp_collapsingmergetree;

-- 執行磁區合并操作
optimize table emp_collapsingmergetree;
-- 再次查詢
select * from emp_collaspingmergetree;

資料折疊不是實時的,需要后臺進行 Compaction 操作,用戶也可以使用手 動合并命令,但是效率會很低,一般不推薦在生產環境中使用,
當進行匯總資料操作時,可以通過改變查詢方式,來過濾掉被洗掉的資料,
SELECT
emp_id,
name,
sum(salary * sign)
FROM emp_collapsingmergetree
GROUP BY emp_id,
name HAVING sum(sign) > 0;

只有相同磁區內的資料才有可能被折疊,其實,當我們修改或洗掉資料時,這些被修改的資料通常是在一個磁區內的,所以不會產生影響,
資料寫入順序,值得注意的是: CollapsingMergeTree對于寫入資料的順序有
著嚴格要求,否則導致無法正常折疊,
如果資料的寫入程式是單執行緒執行的,則能夠較好地控制寫入順序;如果需要處理的資料量很大,資料的寫入程式通常是多執行緒執行的,那么此時就不能保障資料的寫入順序了,在這種情況下,CollapsingMergeTree 的作業機制就會出現問題,但是可以通VersionedCollapsingMergeTree的表引擎得到解決,
3.4.6 VersionedCollapsingMergeTree
上面提到CollapsingMergeTree表引擎對于資料寫入亂序的情況下,不能夠實作資料折疊的效果,VersionedCollapsingMergeTree 表引擎的作用與CollapsingMergeTree完全相同,它們的不同之處在于,VersionedCollapsingMergeTree對資料的寫入順序沒有要求,在同一個磁區內,任意順序的資料都能夠完成折疊操作,
VersionedCollapsingMergeTree使用version列來實作亂序情況下的資料折疊,

該引擎除了需要指定一個sign標識之外,還需要指定--個UInt8型別的 version 版本號,
示例:
CREATE TABLE emp_versioned (
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資',
sign Int8,
version Int8
)ENGINE=VersionedCollapsingMergeTree(sign, version)
PARTITION BY work_place
ORDER BY (emp_id,name)
;
-- 先插入需要被洗掉的資料,即 sign=-1 的資料
INSERT INTO emp_versioned VALUES (1,'tom','上海',25,'技術部',20000,-1,1);
-- 再插入 sign=1 的資料
INSERT INTO emp_versioned VALUES (1,'tom','上海',25,'技術部',20000,1,1);
-- 在插入一個新版本資料
INSERT INTO emp_versioned VALUES (1,'tom','上海',25,'技術部',30000,1,2);

-- 獲取正確查詢結果
SELECT
emp_id,
name,
sum(salary * sign)
FROM emp_versioned
GROUP BY emp_id, name
HAVING sum(sign) > 0;

-- 手動合并
optimize table emp_versioned;
-- 再次查詢
select * from emp_versioned;

雖然在插入資料亂序的情況下,依然能夠實作折疊的效果,之所以能夠達到這種效果,是因為在定義version欄位之后,VersionedCollapsingMergeTree 會自動將version作為排序條件并增加到ORDERBY的末端,就上述的例子而言,最終的排序欄位為ORDER BY emp_ id,name, version desc,
3.5 外部集成表引擎
ClickHouse提供了許多與外部系統集成的方法,包括一些表引擎,這些表引擎與其他型別的表引擎類似,可以用于將外部資料匯入到ClickHouse中,或者在ClickHouse中直接操作外部資料源,
例如直接讀取HDFS的檔案或者MySQL資料庫的表,這些表引擎只負責元資料管理和資料查詢,而它們自身通常并不負責資料的寫入,資料檔案直接由外部系統提供,目前ClickHouse提供了下面的外部集成表引擎:
JDBC:通過指定jdbc連接讀取資料源;
MySQL:將MySQL作為資料存盤,直接查詢其資料;
HDFS:直接讀取HDFS.上的特定格式的資料檔案;
Kafka:將Kafka資料匯入ClickHouse,
3.5.1 HDFS
格式:
ENGINE = HDFS(URI, format)
URI: HDFS 檔案路徑
format:檔案格式,比如CSV、JSON、TSV等,
示例:
CREATE TABLE hdfs_engine_table(
emp_id UInt16 COMMENT '員工 id',
name String COMMENT '員工姓名',
work_place String COMMENT '作業地點',
age UInt8 COMMENT '員工年齡',
depart String COMMENT '部門',
salary Decimal32(2) COMMENT '工資'
) ENGINE=HDFS('hdfs://master:9000/ch/hdfs_engine_table', 'CSV');
-- 寫入資料
INSERT INTO hdfs_engine_table VALUES (1,'tom','上海',25,'技術部',20000),(2,'jack','上海',26,'人事部',10000);
注意:
ch 目錄要存在
如果提示沒權限:hdfs dfs -chmod -R 777 /
可以看出,這種方式與使用Hive類似,我們直接可以將HDFS對應的檔案映射成ClickHouse中的一張表,這樣就可以使用SQL操作HDFS上的檔案了,
值得注意的是:ClickHouse并不能夠洗掉HDFS上的資料,當我們在ClickHouse客戶端中洗掉了對應的表,只是洗掉了表結構,HDFS.上的檔案并沒有被洗掉,這一點跟Hive的外部表十分相似,

3.5.2 MySQL
格式:
CREATE TABLE mysql_engine_table(
id Int32,
name String
) ENGINE = MySQL( 'master:3306','test', 'student', 'root', '123456');
-- 查詢資料
SELECT * FROM mysql_engine_table;

對于MySQL表引擎,不支持UPDATE和DELETE操作,
3.5.3 Kafka
格式:
示例:
CREATE TABLE kafka_table (
uid UInt64,
phone UInt64,
addr String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'master:9092',
kafka_topic_list = 'ck_topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow' ;
生產者:
kafka-console-producer.sh --broker-list master:9092 --topic ck_topic
資料:
{"uid":"1000166111","phone":"17703771999","addr":" 河南省 南陽"}{"uid":"1000432103","phone":"15388889881","addr":" 云南省 昆明"}{"uid":"1000473355","phone":"15388889557","addr":" 云南省 昆明"}{"uid":"1000555472","phone":"18083815777","addr":" 云南省 昆明"}{"uid":"1000585644","phone":"15377892222","addr":" 廣東省 中山"}{"uid":"1000774061","phone":"18026666666","addr":" 廣東省 惠州"}{"uid":"1001024965","phone":"18168526111","addr":" 江蘇省 蘇州"}{"uid":"1001283200","phone":"15310952123","addr":" 重慶 "}{"uid":"1001523180","phone":"15321168157","addr":" 北京 "}

當我們一旦查詢完畢之后,ClickHouse 會洗掉表內的資料,其實Kafka表引擎只是一個資料管道,我們可以通過物化視圖的方式訪問Kafka中的資料,
首先創建--張Kafka表引擎的表,用于從Kafka中讀取資料;
然后再創建--張普通表引擎的表,比如MergeTree,面向終端用戶使用;
最后創建物化視圖,用于將Kafka引擎表實時同步到終端用戶所使用的表中,
示例:
-- 創建 Kafka 引擎表
CREATE TABLE kafka_table_consumer (
uid UInt64,
phone UInt64,
addr String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'master:9092',
kafka_topic_list = 'ck_topic',
kafka_group_name = 'group1',
kafka_format = 'JSONEachRow' ;
-- 創建一張終端用戶使用的表
CREATE TABLE kafka_table_mergetree (
uid UInt64,
phone UInt64,
addr String
)ENGINE=MergeTree()
ORDER BY uid ;
-- 創建物化視圖,同步資料
CREATE MATERIALIZED VIEW consumer TO kafka_table_mergetree AS
SELECT uid,phone,addr
FROM kafka_table_consumer ;
-- 查詢,多次查詢,已經被查詢的資料依然會被輸出
select * from kafka_table_mergetree;

3.6 其他引擎
3.6.1 Memory
記憶體引擎,資料以未壓縮的原始形式直接保存在記憶體當中,服務器重啟資料就會消失,
讀寫操作不會相互阻塞,不支持索引,簡單查詢下有非常非常高的性能表現(超過10G/s )
一般用到它的地方不多,除了用來測驗,就是在需要非常高的性能,同時資料量又不太大(上限大概1億行)的場景,
使用方式和TinyLog一致,
3.6.2 Distributed
Distributed表引擎是分布式表的代名詞,它自身不存盤任何資料,資料都分散存盤在某一個分片上,能夠自動路由資料至集群中的各個節點,所以Distributed表引擎需要和其他資料表引擎一起協同作業,
所以,一張分布式表底層會對應多個本地分片資料表,由具體的分片表存盤資料,分布式表與分片表是一對多的關系,
Distributed(cluster_ name, database_ name, table_ _name[, sharding_ key])
分布式引擎引數:服務器組態檔中的集群名,遠程資料庫名,遠程表名,資料分片鍵(可選),
1. 創建分布式表
默認情況下,CREATE、DROP、ALTER、RENAME操作僅僅在當前執行該命令的server.上生效,在集群環境下,可以使用ON CLUSTER陳述句,這樣就可以在整個集群發揮作用,
創建一張分布式表:
CREATE TABLE IF NOT EXISTS user_cluster ON CLUSTER cluster_3shards_1replicas (
id Int32,
name String
)ENGINE = Distributed(cluster_3shards_1replicas, default, user_local,id);
Distributed表引擎的定義形式如下所示:
Distributed(cluster_name, database_name, table_name[, sharding_key])各個引數的含義分別如下:cluster_name :集群名稱,與集群配置中的自定義名稱相對應,database_name :資料庫名稱table_name :表名稱sharding_key :可選的,用于分片的 key 值,在資料寫入的程序中,分布式表會依據分片 key 的規則,將資料分布到各個節點的本地表,
使用了ON CLUSTER分布式DDL,這意味著在集群的每個分片節點上,都會創建一張Distributed表,這樣便可以從其中任意一- 端發起對所有分片的讀、寫請求,在每臺機器上查看表,發現每臺機器上都存在一張剛剛創建好的表,
2. 創建本地表
在每臺機器_上分別創建一張本地表:
CREATE TABLE IF NOT EXISTS user_local (
id Int32,
name String
)ENGINE = MergeTree()
PARTITION BY id
ORDER BY id
PRIMARY KEY id;
-- 查詢本地表
select * from user_local;
先在一臺機器.上對user_ local 表進行插入資料,然后再查詢user_ cluster 表,

再向user_ cluster 中插入一些資料,觀察user_ local 表資料變化,可以發現資料被分散存盤到了其他節點上了,
INSERT INTO user_cluster VALUES(3,'lilei'),(4,'lihua'); master
-- 在master節點上查詢
select * from user_local;

-- 在 slave1節點上查詢
select * from user_local

四、Flink 整合ClickHouse
4.1初步整合實作
package cn.kgc.sink
import java.sql.PreparedStatement
import org.apache.flink.connector.jdbc._
import org.apache.flink.streaming.api.scala._
object CKSink {
def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.createLocalEnvironment();
val datastream = env.fromElements((7,"lele"),(8,"xixi"))
val sql = "insert into t2 values(?,?)"
datastream.addSink(JdbcSink
.sink[(Int,String)](
sql,
new CkSinkBuilder,
new JdbcExecutionOptions.Builder().withBatchSize(5).build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl("jdbc:clickhouse://master:8123").withDriverName("ru.yandex.clickhouse.ClickHouseDriver") .withUsername("default") .build()
))
env.execute()
}
}
class CkSinkBuilder extends JdbcStatementBuilder[(Int, String)] {
override def accept(t: PreparedStatement, u: (Int, String)): Unit = {
t.setInt(1, u._1)
t.setString(2, u._2)
}
}
4.2 專案資料寫入 ClickHouse
package cn.kgc.stock
import java.sql.PreparedStatement
import java.text.SimpleDateFormat
import java.util.{Date, Properties}
import com.alibaba.fastjson.{JSON, JSONArray}
import org.apache.flink.api.common.functions.RichFlatMapFunction
import org.apache.flink.api.common.restartstrategy.RestartStrategies
import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.api.common.time.Time import org.apache.flink.connector.jdbc.{JdbcConnectionOptions, JdbcExecutionOptions, JdbcSink, JdbcStatementBuilder}
import org.apache.flink.runtime.state.hashmap.HashMapStateBackend
import org.apache.flink.streaming.api.CheckpointingMode
import org.apache.flink.streaming.api.environment.CheckpointConfig.ExternalizedCheckpoin tCleanup
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer
import org.apache.flink.util.Collector
import org.apache.kafka.clients.consumer.ConsumerConfig
object Kafka2Flink2ClickHouse {
def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setParallelism(1)
val hashMapStateBackend = new HashMapStateBackend()
env.setStateBackend(new HashMapStateBackend())
try {
//env.getCheckpointConfig.setCheckpointStorage("file:///D://abc//ckp")
env.getCheckpointConfig
.setCheckpointStorage("hdfs://master:9000/flink/checkpoin t")
} catch {
case e => e.printStackTrace()
}
env.enableCheckpointing(1000,CheckpointingMode.EXACTLY_ONCE)
env.getCheckpointConfig.setMinPauseBetweenCheckpoints(500)
env.getCheckpointConfig.setCheckpointTimeout(60000)
env.getCheckpointConfig.setTolerableCheckpointFailureNumber(1)
env.getCheckpointConfig.enableExternalizedCheckpoints(ExternalizedCheckpointClea nup.RETAIN_ON_CANCELLATION)
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3,Time.milliseconds(600)))
val props = new Properties()
props.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "master:9092")
props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "group-2")
props.setProperty(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "2000")
props.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true")
props.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")
val kafkaConsumer = new FlinkKafkaConsumer[String]("indexTimeLine",new SimpleStringSchema(),props)
kafkaConsumer.setCommitOffsetsOnCheckpoints(true)
val inputStream = env.addSource(kafkaConsumer)
val resultStream: DataStream[StockIndexMinuteResult2] = inputStream.flatMap(new MyFlatMapAnalysisJson2)
resultStream.print()
val sql = "insert into t5 values(?,?,?)"
resultStream.addSink(JdbcSink
.sink[StockIndexMinuteResult2]( sql, new CkSinkBuilder, new JdbcExecutionOptions.Builder().withBatchSize(5).build(),
new JdbcConnectionOptions
.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:clickhouse://master:8123")
.withDriverName("ru.yandex.clickhouse.ClickHouseDriver")
.withUsername("default")
.build()
)
)
env.execute()
}
}
class CkSinkBuilder extends JdbcStatementBuilder[StockIndexMinuteResult2] {
override def accept(t: PreparedStatement, u: StockIndexMinuteResult2): Unit = {
t.setString(1,u.name)
t.setString(2,u.dt)
t.setDouble(3,u.nowPrice)
}
}
class MyFlatMapAnalysisJson2 extends RichFlatMapFunction[String,StockIndexMinuteResult2]{
override def flatMap(line: String, out: Collector[StockIndexMinuteResult2]): Unit = {
val jSONObject = JSON.parseObject(line)
val jSONArray = jSONObject.getJSONObject("showapi_res_body").getJSONArray("dataList")
val jSONObject2 = JSON.parseObject(jSONArray.getString(0))
val dt = jSONObject2.get("date").toString
val array: JSONArray = jSONObject2.getJSONArray("minuteList")
val code = jSONObject.getJSONObject("showapi_res_body").getString("code")
val name = jSONObject.getJSONObject("showapi_res_body").getString("name")
val market = jSONObject.getJSONObject("showapi_res_body").getString("market")
var i:Int = 0
while (i < array.size()){
val time = array.getJSONObject(i).getString("time")
val avgPrice = array.getJSONObject(i).getString("avgPrice").toDouble
val nowPrice = array.getJSONObject(i).getString("nowPrice").toDouble
out.collect( StockIndexMinuteResult2(code,name,market,tranTimeToString(tranTimeToLong(dt+ti me)),avgPrice,nowPrice))
i += 1
}
}
def tranTimeToLong(tm:String) :Long={
val fm = new SimpleDateFormat("yyyyMMddHHmm")
val dt = fm.parse(tm)
val aa = fm.format(dt)
val tim: Long = (dt.getTime()+"").toLong tim
}
def tranTimeToString(tm:Long) :String={
val fm = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
val time = fm.format(new Date(tm)) time
}
}
case class StockIndexMinuteResult2(codeId:String,name:String,market:String,dt:String,avgPrice :Double,nowPrice:Double)
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/298622.html
標籤:其他
