生產者創建訊息,在其他基于發布與訂閱的訊息系統中,生產者可能被稱為發布者 或 寫入者,
一般情況下,一個訊息會被發布到一個特定的主題上,生產者在默認情況下把訊息均衡地分布到主題的所有磁區上,而并不關心特定訊息會被寫到哪個磁區,不過,在某些情況下,生產者會把訊息直接寫到指定的磁區,這通常是通過訊息鍵和磁區器來實作的,磁區器為鍵生成一個散列值,并將其映射到指定的磁區上,這樣可以保證包含同一個鍵的訊息會被寫到同一個磁區上,生產者也可以使用自定義的磁區器,根據不同的業務規則將訊息映射到磁區,
生產者發送訊息的方式
生產者發送訊息主要有 2 種方式:同步發送訊息、異步發送訊息
同步發送訊息
同步發送訊息:我們呼叫 KafkaProducer 的 send() 方法發送訊息,send() 方法會回傳一個包含 RecordMetadata 的 Future 物件,然后呼叫 Future 的 get() 方法等待 Kafka 回應,通過 Kafka 的回應,我們就可以知道訊息是否發送成功,
- 如果服務器回傳錯誤,Future 的 get() 方法會拋出例外,
- 如果沒有發生錯誤,我們會得到一個 RecordMetadata 物件,這個物件包含訊息的目標主題、磁區資訊和訊息的偏移量等資訊,
我們呼叫 KafkaProducer 的 send() 方法發送 ProducerRecord 物件,訊息先是被放進緩沖區,然后使用單獨的執行緒將訊息發送到服務器端,
例外處理
如果在發送資料之前或者在發送程序中發生了任何錯誤,比如 broker 回傳了一個不允許重發訊息的例外或者已經超過了重發的次數,那么就會拋出例外,在發送訊息之前,生產者也是有可能發生例外的,這些例外有可能是 SerializationException(說明序列化訊息失敗)、BufferExhaustedException 或 TimeoutException(說明緩沖區已滿),又或者是 InterruptException(說明發送執行緒被中斷),
KafkaProducer 一般會發生兩類錯誤,
- 其中一類是可重試錯誤,這類錯誤可以通過重發訊息來解決,比如對于連接錯誤,可以通過再次建立連接來解決,“無主(no leader)”錯誤則可以通過重新為磁區選舉首領來解決,KafkaProducer 可以被配置成自動重試,如果在多次重試后仍無法解決問題,應用程式會收到一個重試例外,
- 另一類錯誤無法通過重試解決,比如“訊息太大”例外,對于這類錯誤,KafkaProducer 不會進行任何重試,直接拋出例外,
public void send(String topic, String key, String val) {
ProducerRecord<String, String> producerRecord = new ProducerRecord<>(topic, key, val);
try {
ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(producerRecord);
SendResult<String, String> sendResult = future.get();
} catch (Exception e) {
e.printStackTrace();
}
}
異步發送訊息
異步發送訊息:我們呼叫 KafkaProducer 的 send() 方法,并指定一個回呼方法,在服務器回傳回應時呼叫該方法,
大多數時候,我們并不需要等待回應,不過在遇到訊息發送失敗時,我們需要拋出例外、記錄錯誤日志,或者把訊息寫入“錯誤訊息”檔案以便日后分析,為了在異步發送訊息的同時能夠對例外情況進行處理,生產者提供了回呼支持,
為了使用回呼,需要一個實作了 org.apache.kafka.clients.producer.Callback 介面的類,這個介面只有一個 onCompletion() 方法,如果 Kafka 回傳一個錯誤,onCompletion() 方法會拋出一個非空例外,通過 onCompletion() 方法拋出的例外,我們可以對發送失敗的訊息進行處理,一般情況下,因為生產者會自動進行重試,所以就沒必要在代碼邏輯里處理那些可重試的錯誤,你只需要處理那些不可重試的錯誤或重試次數超出上限的情況,
private class DemoProducerCallback implements Callback {
@Override
public void onCompletion(RecordMetadata recordMetadata, Exception e) {
if (e != null) {
e.printStackTrace();
}
}
}
ProducerRecord<String, String> record = new ProducerRecord<>("CustomerCountry", "Biomedical Materials", "USA");
producer.send(record, new DemoProducerCallback());
磁區器
介紹磁區
ProducerRecord 物件包含目標主題、訊息鍵和值(訊息),
- 如果訊息鍵為 null,并且使用了默認的 DefaultPartitioner 磁區器,那么磁區器使用粘性磁區策略(UniformSticky),會隨機選擇一個磁區,并盡可能一直使用該磁區,等到該磁區的 batch 已滿或者已完成,Kafka 再隨機一個磁區進行使用(保證和上一次的磁區不同),
- 如果訊息鍵不為 null,并且使用了默認的 DefaultPartitioner 磁區器,那么磁區器會對訊息鍵進行散列(使用 Kafka 自己的散列演算法,即使升級 Java 版本,散列值也不會發生變化),然后根據散列值把訊息映射到特定的磁區上(散列值 與 主題的磁區數進行取余得到 partition 值),
這里的關鍵之處在于,同一個鍵總是被映射到同一個磁區上,所以在進行映射時,我們會使用主題的所有磁區,而不僅僅是可用的磁區,這也意味著,如果寫入資料的磁區是不可用的,那么就會發生錯誤,
只有在不改變主題磁區數量的情況下,鍵與磁區之間的映射才能保持不變,一旦主題增加了新的磁區,那么鍵與磁區之間的映射關系就改變了,如果要使用鍵來映射磁區,那么最好在創建主題的時候就把磁區規劃好,而且永遠不要增加新磁區,
自定義磁區策略
生產者可以使用自定義的磁區器,根據不同的業務規則將訊息映射到磁區,
通過磁區器實作自定義磁區策略的步驟:
- 定義一個類,該類實作 Partitioner 介面(磁區器)
- 配置生產者(KafkaProducer),讓生產者發送訊息時使用自定義的磁區器:properties.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, MyPartitioner.class.getName());
public class MyPartitioner implements Partitioner {
/**
* 回傳資訊對應的磁區
*
* @param topic 主題
* @param key 訊息的 key
* @param keyBytes 訊息的 key 序列化后的位元組陣列
* @param value 訊息的 value
* @param valueBytes 訊息的 value 序列化后的位元組陣列
* @param cluster 集群元資料可以查看磁區資訊
* @return
*/
@Override
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
if ((keyBytes == null) || (!(key instanceof String))) {
throw new InvalidRecordException("We expect all messages to have String type as key");
}
// 實作自己的磁區策略
// 回傳資料寫入的磁區號
return 0;
}
// 關閉資源
@Override
public void close() {
}
// 配置方法
@Override
public void configure(Map<String, ?> configs) {
}
}
參考資料
《Kafka 權威指南》第 3 章:Kafka 生產者——向 Kafka 寫入資料
本文來自博客園,作者:真正的飛魚,轉載請注明原文鏈接:https://www.cnblogs.com/feiyu2/p/17250979.html
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/547981.html
標籤:大數據
