kafka初學
一、介紹
- Kafka是是一個分布式、支持磁區的(partition)、多副本的(replica),基于zookeeper協調的分布式訊息系統
- 它的最大的特性就是可以實時的處理大量資料以滿足各種需求場景:
- 比如基于hadoop的批處理系統
- 低延遲的實時系統
- Storm/Spark流式處理引擎
- web/nginx日志
- 訪問日志
- 訊息服務等等
- Kafka用scala語言撰寫
1. 應用場景
- 日志收集:一個公司可以用Kafka收集各種服務的log,通過kafka以統一介面服務的方式開放給各種consumer,例如hadoop、Hbase、Solr等,
- 訊息系統:解耦和生產者和消費者、快取訊息等,
- 用戶活動跟蹤:Kafka經常被用來記錄web用戶或者app用戶的各種活動,如瀏覽網頁、搜索、點擊等活動,這些活動資訊被各個服務器發布到kafka的topic中,然后訂閱者通過訂閱這些topic來做實時的監控分析,或者裝載到hadoop、資料倉庫中做離線分析和挖掘,
- 運營指標:Kafka也經常用來記錄運營監控資料,包括收集各種分布式應用的資料,生產各種操作的集中反饋,比如報警和報告,
2. Kafka基本概念
- kafka是一個分布式的,磁區的訊息服務
- kafka提供一個訊息系統應該具備的功能,但是確有著獨特的設計,可以這樣來說,Kafka借鑒了JMS規范的思想,但是確并沒有完全遵循JMS規范,
2.1 基礎的訊息(Message)相關術語
| 名稱 | 解釋 |
|---|---|
| Broker | 訊息中間件處理節點,一個Kafka節點就是一個broker,一個或者多個Broker可以組成一個Kafka集群 |
| Topic | Kafka根據topic對訊息進行歸類,發布到Kafka集群的每條訊息都需要指定一個topic |
| Producer | 訊息生產者,向Broker發送訊息的客戶端 |
| Consumer | 訊息消費者,從Broker讀取訊息的客戶端 |
| ConsumerGroup | 每個Consumer屬于一個特定的Consumer Group,一條訊息可以被多個不同的Consumer Group消費,但是一個Consumer Group中只能有一個Consumer能夠消費該訊息 |
| Partition | 物理上的概念,一個topic可以分為多個partition,每個partition內部訊息是有序的 |
從一個較高的層面上來看,producer通過網路發送訊息到Kafka集群,然后consumer來進行消費:
- 服務端(brokers)和客戶端(producer、consumer)之間通信通過TCP協議來完成,
二、kafka基本使用
1. 安裝前的環境準備
- 安裝jdk
- 安裝zk
- 官網下載kafka的壓縮包:http://kafka.apache.org/downloads
- 解壓縮至如下路徑
/usr/local/kafka/
- 修改組態檔:/usr/local/kafka/kafka2.11-2.4/config/server.properties
#broker.id屬性在kafka集群中必須要是唯一
broker.id=0
#kafka部署的機器ip和提供服務的埠號
listeners=PLAINTEXT://192.168.65.60:9092
#kafka的訊息存盤檔案
log.dir=/usr/local/data/kafka-logs
#kafka連接zookeeper的地址
zookeeper.connect=192.168.65.60:2181
2.啟動kafka服務器
進入到bin目錄下,使用命令來啟動
./kafka-server-start.sh -daemon ../config/server.properties
驗證是否啟動成功:
進入到zk中的節點看id是0的broker有沒有存在(上線)
ls /brokers/ids/
server.properties核心配置詳解:
| Property | Default | Description |
|---|---|---|
| broker.id | 0 | 每個broker都可以用一個唯一的非負整數id進行標識;這個id可以作為broker的“名字”,你可以選擇任意你喜歡的數字作為id,只要id是唯一的即可, |
| log.dirs | /tmp/kafka-logs | kafka存放資料的路徑,這個路徑并不是唯一的,可以是多個,路徑之間只需要使用逗號分隔即可;每當創建新partition時,都會選擇在包含最少partitions的路徑下進行, |
| listeners | PLAINTEXT://192.168.65.60:9092 | server接受客戶端連接的埠,ip配置kafka本機ip即可 |
| zookeeper.connect | localhost:2181 | zooKeeper連接字串的格式為:hostname:port,此處hostname和port分別是ZooKeeper集群中某個節點的host和port;zookeeper如果是集群,連接方式為 hostname1:port1, hostname2:port2, hostname3:port3 |
| log.retention.hours | 168 | 每個日志檔案洗掉之前保存的時間,默認資料保存時間對所有topic都一樣, |
| num.partitions | 1 | 創建topic的默認磁區數 |
| default.replication.factor | 1 | 自動創建topic的默認副本數量,建議設定為大于等于2 |
| min.insync.replicas | 1 | 當producer設定acks為-1時,min.insync.replicas指定replicas的最小數目(必須確認每一個repica的寫資料都是成功的),如果這個數目沒有達到,producer發送訊息會產生例外 |
| delete.topic.enable | false | 是否允許洗掉主題 |
3.創建主題topic
topic是什么概念?topic可以實作訊息的分類,不同消費者訂閱不同的topic,
[外鏈圖片轉存失敗,源站可能有防盜鏈機制,建議將圖片保存下來直接上傳(img-6028PZRb-1646227957372)(img/截屏2021-07-08 下午2.41.33.png)]
執行以下命令創建名為“test”的topic,這個topic只有一個partition,并且備份因子也設定為1:
./kafka-topics.sh --create --zookeeper 172.16.253.35:2181 --replication-factor 1 --partitions 1 --topic test
查看當前kafka內有哪些topic
./kafka-topics.sh --list --zookeeper 172.16.253.35:2181
4.發送訊息
kafka自帶了一個producer命令客戶端,可以從本地檔案中讀取內容,或者我們也可以以命令列中直接輸入內容,并將這些內容以訊息的形式發送到kafka集群中,在默認情況下,每一個行會被當做成一個獨立的訊息,使用kafka的發送訊息的客戶端,指定發送到的kafka服務器地址和topic
./kafka-console-producer.sh --broker-list 172.16.253.21:9092 --topic test
5.消費訊息
對于consumer,kafka同樣也攜帶了一個命令列客戶端,會將獲取到內容在命令中進行輸出,默認是消費最新的訊息,使用kafka的消費者訊息的客戶端,從指定kafka服務器的指定topic中消費訊息
- 方式一:從最后一條訊息的偏移量+1開始消費
./kafka-console-consumer.sh --bootstrap-server 172.16.253.21:9092 --topic test
- 方式二:從頭開始消費
./kafka-console-consumer.sh --bootstrap-server 172.16.253.21:9092 --from-beginning --topic test
幾個注意點:
- 訊息會被存盤
- 訊息是順序存盤
- 訊息是有偏移量的
- 消費時可以指明偏移量進行消費
三、Kafka中的關鍵細節
1. 訊息的順序存盤
- 訊息的生產者會把訊息發送到broker中,broker會存盤訊息,訊息是按照順序進行存盤的,
- 訊息消費者在消費訊息的時候也是按照順序消費的,消費訊息時可以從默認位置(最后一條訊息的下一個偏移量)開始消費,也可以指定某個位置開始消費

2. 單播訊息的實作
- 單播訊息:一個消費組里,最多只有一個消費者能消費到某一個topic中的訊息,除非次消費者掛掉,否則下一個消費者無法消費到此topic中的訊息
./kafka-console-consumer.sh --bootstrap-server 172.16.253.21:9092 --consumer-property group.id=testGroup --topic test
3. 多播訊息的實作
- 不同消費者組中的某一個消費者可以同時收到訊息
- 在一些業務場景中需要讓一條訊息被多個消費者消費,那么就可以使用多播模式,
./kafka-console-consumer.sh --bootstrap-server 10.31.167.10:9092 --consumer-property group.id=testGroup1 --topic test
./kafka-console-consumer.sh --bootstrap-server 10.31.167.10:9092 --consumer-property group.id=testGroup2 --topic test

4. 查看訊息組及資訊
- Currennt-offset: 當前消費組的已消費偏移量
- Log-end-offset: 主題對應磁區訊息的結束偏移量(HW)
- Lag: 當前消費組未消費的訊息數

# 查看當前主題下有哪些消費組
./kafka-consumer-groups.sh --bootstrap-server 172.16.253.21:9092 --list
# 查看消費組中的具體資訊:比如當前偏移量、最后一條訊息的偏移量、堆積的訊息數量
./kafka-consumer-groups.sh --bootstrap-server 172.16.253.21:9092 --describe --group testGroup

四、主題、磁區的概念
1. 主題Topic
- 主題Topic可以理解成是一個類別的名稱
2. 磁區partition
- 一個主題中的訊息量是非常大的,因此可以通過磁區的設定,來分布式存盤這些訊息,比如一個topic創建了3個磁區,那么topic中的訊息就會分別存放在這三個磁區中,
- 為一個主題創建多個磁區
./kafka-topics.sh --create --zookeeper 172.16.253.81:2181 --partitions 2 --replication-factor 1 --topic test1
- 可以通過這樣的命令查看topic的磁區資訊
./kafka-topics.sh --describe --zookeeper 172.16.253.81:2181 --topic test1
- 磁區的作用:
- 可以分布式存盤
- 可以并行寫
實際上是存在data/kafka-logs/test-0 和 test-1中的0000000.log檔案中
小細節:
-
定期將??消費磁區的offset提交給kafka內部topic:__consumer_offsets,提交過去的時候,key是consumerGroupId+topic+磁區號,value就是當前offset的值,kafka會定期清理topic?的訊息,最后就保留最新的那條資料
因為__consumer_offsets可能會接收?并發的請求,kafka默認給其分配50個磁區(可以通過offsets.topic.num.partitions設定),這樣可以通過加機器的?式抗?并發,
通過如下公式可以選出consumer消費的offset要提交到__consumer_offsets的哪個磁區
公式:hash(consumerGroupId) % __consumer_offsets主題的磁區數
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/437058.html
標籤:其他
上一篇:Spark standalone模式在多用戶環境下保存結果報錯 java.io.ioexception: mkdirs failed to create file
