產生訊息
生產者是一個將訊息串流發佈至 Kafka 主題的應用程式。 此資訊主要討論作為 Apache Kafka 專案一部分的 Java 程式設計介面。 這些概念也套用於其他語言,但名稱有時會稍有不同。
在程式設計介面中,訊息稱為記錄。 例如,Java 類別 org.apache.kafka.clients.producer.ProducerRecord 用來表示來自生產者 API 觀點的訊息。 術語_記錄_和_訊息_可以互換使用,但本質上記錄用於表示訊息。
當生產者連接至 Kafka 時,即會起始引導連線。 此連線可連接至叢集裡的任何伺服器。 生產者會要求與其想要發佈至其中的主題如需的分割區及領導權資訊。 然後,生產者與分區領導者建立另一個連結並可以發布訊息。 當您的生產者連接至 Kafka 叢集時,即會在內部自動發生這些動作。
為了確保可用性,Kafka 分配管理系統會抄寫訊息,因此如果有一個分配管理系統無法使用,其他分配管理系統仍可以接收來自生產者的訊息,並將它們傳送給消費者。Event Streams 使用抄寫因數 3,表示每一則訊息都儲存在三個分配管理系統上。 當訊息傳送至分割區領導者時,消費者無法立即取用該訊息。 領導者會將訊息記錄附加至分割區,為其指派該分割區的下一個偏移數。 在同步副本的所有追蹤者複製記錄並確認將記錄寫入其副本後,該記錄現在已提交並可供使用者使用。
每個訊息都表示為一筆記錄,由兩部分組成:鍵和值。 索引鍵通常用於有關訊息的資料,值則是訊息內文。 由於Kafka生態系統中的許多工具(例如連接其他系統的連接器)僅使用值而忽略鍵,因此最好將所有訊息資料放入值中並使用鍵進行分區或日誌壓縮。 不要依賴從Kafka讀取的所有內容來使用金鑰。
許多其他傳訊系統也有一種攜帶其他資訊與訊息的方式。 Kafka 0.11版本為此引進了記錄標頭。
您可能會發現在 Event Streams中閱讀此資訊以及 耗用訊息 可能會很有用。
配置設定
生產者有許多配置設定。 您可以控制生產者的各個方面,包括批次、重試和訊息確認。 以下是最重要的部分:
| 名稱 | 說明 | 有效值 | 預設值 |
|---|---|---|---|
| key.serializer | 此類別用來序列化索引鍵。 | 實作 Serializer 介面的Java類,例如org.apache.kafka.common.serialization.StringSerializer。 | 無預設值 - 您必須指定一個值。 |
| value.serializer | 此類別用來序列化值。 | 實作 Serializer 介面的Java類,例如org.apache.kafka.common.serialization.StringSerializer。 | 無預設值 - 您必須指定一個值。 |
| acks | 用來確認每一則已發佈訊息所需的伺服器數目。 這會控制生產者所需的延續性保證。 | 0、1、all(或 -1) | all (Kafka 3.0 以及更新版本) 1 (早於 Kafka 3.0) |
| retries | 傳送發生錯誤時,用戶端重送訊息的次數。 | 0,... | 0 |
| max.block.ms | 傳送或 meta 資料要求可以阻擋等待的毫秒數。 | 0,... | 60000(1 分鐘) |
| max.in.flight.requests.per.connection | 客戶端在阻止進一步請求之前在連線上發送的未確認請求的最大數量。 | 1,... | 5 |
| request.timeout.ms | 生產者等待要求回應的最長時間量。 如果在逾時之前未收到回應,則請求將重試,如果重試次數已用完,則請求失敗。 | 0,... | 30000(30 秒) |
還有更多配置設定可用,但請確保在嘗試之前仔細閱讀 Apache Kafka文件。
分割
在 Kafka 中,分割區是可調整性的單元。 因此,分區是提高吞吐量的有效方法,因為它允許主題資料在多個並行流中流動。
當生產者在主題上發佈訊息時,生產者可以選擇要使用的分割區。 如果排序很重要,請注意分區是記錄的有序序列,但主題包含一個或多個分區。 如果您要依序遞送一組訊息,請確定它們全部都在同一個分割區中。 達成此目的的最直接方式為對所有這些訊息提供相同的索引鍵。
在生產者發佈訊息時,可以明確指定分割區號碼。 這提供直接控制,但會讓生產者程式碼更複雜,因為它要負責管理分割區選擇。 如需相關資訊,請參閱呼叫 Producer.partitionsFor 的方法。 例如,該呼叫是針對 Kafka版本2.2.0 進行描述的。
如果生產者未指定分割區號碼,則由分割程式進行分割區選擇。 Kafka 生產者的內建預設分割程式的運作方式如下:
- 如果記錄沒有索引鍵,請依循環方式選取分割區。
- 如果記錄具有索引鍵,請藉由計算索引鍵的雜湊值來選取分割區。 這具有為所有含有相同索引鍵的訊息選取相同分割區的效果。
您也可以撰寫自己的自訂分割程式。 自訂分割程式可以選擇任何方法來將記錄指派給分割區。 例如,只使用索引鍵中的資訊子集或使用應用程式特有的 ID。
訊息排序
Kafka 通常會依據生產者傳送的順序來撰寫訊息。 但是,在某些情況下,重試可能會導致訊息重複或重新排序。 如果您希望按順序發送一系列訊息,那麼確保它們都寫入同一個分區非常重要,因為這是保證訊息排序的唯一方法。
生產者也能夠自動重試發送訊息。 啟用此重試功能是個好主意,因為另一種選擇是您的應用程式程式碼必須自行執行任何重試。 Kafka 中的分批處理及自動重試的組合可具有複製訊息並重新予以排序的效果。
例如,如果您在某個主題上發布三個訊息<M1, M2, M3>的序列。 這些記錄可能都適合同一批次,因此它們都會一起發送給分區領導者。 然後,領導者會將它們寫入分割區,並將它們作為個別記錄進行抄寫。 如果發生故障,可能會將 M1 和 M2 新增至分割區,但不會將 M3 新增至分割區。 生產者未收到確認通知,因此它會重試傳送 <M1、M2、M3>。新的主導器會將 M1、M2 及 M3 寫入分割區,該分割區現在包含 <M1、M2、M1、M2、M3>,其中複製的 M1 遵循原始 M2。 如果您將每一個分配管理系統的進行中要求數目限制為一個,則可以防止此重新排序。 您可能仍會發現單一記錄是重複的,例如 <M1、M2、M2、M3>,但永遠不會出現不正常的順序。 在Kafka 0.11或更高版本中,也可以使用冪等生產者功能來防止M2重複。
使用Kafka時,通常的做法是編寫應用程式來處理偶爾出現的訊息重複,因為只有單一請求在運行時對效能有很大的影響。
訊息確認通知
當您發布訊息時,您可以使用 acks 生產者配置來選擇所需的確認等級。 該選項表示傳輸量與可靠性之間的平衡。 存在下列三個層次。
- acks=0(最不可靠)
- 訊息一寫入網路就被視為已發送。 沒有來自分割區領導者的確認通知。 因此,如果分割區領導權發生變更,則訊息可能會遺失。 這種等級的確認速度很快,但在某些情況下可能會遺失訊息。
- acks=1(預設值)
- 一旦分區領導者成功地將其記錄寫入分區,該訊息就會被確認給生產者。 由於確認發生在記錄到達同步副本之前,因此如果領導者發生故障,但追隨者尚未收到訊息,則訊息可能會遺失。 如果分區領導發生變化,舊的領導者會通知生產者,生產者可以處理錯誤並重試將訊息發送給新的領導者。 由於訊息在所有副本確認其接收之前已被確認,因此如果分區領導層發生更改,已確認但尚未完全複製的訊息可能會遺失。
- acks=all(最可靠)
- 當分區領導者成功寫入其記錄並且所有同步副本也執行相同操作時,該訊息將被確認給生產者。 如果分割區領導權發生變更,訊息也不會遺失,前提是至少有一個同步抄本可供使用。
即使您不等待訊息被確認給生產者,訊息仍然只能在已確定時被取用,這表示抄寫到同步抄本已完成。 換言之,從生產者的觀點來看傳送訊息的延遲低於從生產者傳送訊息到接收訊息的消費者所測量的端對端延遲。
如果可能,請避免在發布下一則訊息之前等待訊息確認。 等待防止生產者能夠將訊息分批在一起,並且還可以將訊息發佈的速率降低到低於網路的來回轉換延遲。
分批處理、節流控制及壓縮
出於效率目的,生產者將大量記錄收集在一起以發送到伺服器。 如果啟用壓縮,生產者會壓縮每一個批次,如此即可藉由減少在網路上傳送資料而提高效能。
如果您嘗試發佈訊息的速度比訊息傳送給伺服器的速度快,則生產者會自動將其緩衝至分批處理要求中。 生產者會維護每一個分割區的未傳送記錄的緩衝區。 有時甚至批次也無法達到您想要的速率。
還有另一個影響因素。 為了防止個體生產者或消費者徹底擊敗叢集,Event Streams 會套用傳輸量配額。 計算每一個生產者傳送資料的速率,並節流控制嘗試超出其配額的任何生產者。 藉由稍稍延遲向生產者傳送回應來套用節流控制。 通常,這只是一個自然的約束。
如需傳輸量指引的相關資訊,請參閱 限制和配額。
總之,當訊息發佈時,其記錄會先寫入生產者的緩衝區中。 生產者會在背景中分批處理記錄並將其傳送到伺服器。 然後,伺服器會回應生產者,如果生產者發佈得太快,則可能會套用節流控制延遲。 如果生產者中的緩衝區已滿,生產者的發送呼叫就會延遲,但最終可能會失敗並出現異常。
遞送語意
Kafka 提供下列多種不同的訊息遞送語意:
- 最多一次:訊息可能會遺失並且不會重新傳遞。
- 至少一次:訊息永遠不會遺失,但可能會有重複。
- 恰好一次:訊息永遠不會遺失並且不會重複。
遞送語意由下列設定決定:
acksretriesenable.idempotence
預設情況下,Kafka使用至少一次語義。
要啟用恰好一次語義,您必須使用冪等或事務性生產者。 等冪生產者的啟用方式,是將 enable.idempotence 設為 true,並保證即使重試,也只會將每則訊息的恰好一份寫入 Kafka。 交易式生產者讓您可將資料傳送到多個分割區,以便所有訊息都順利遞送,或是全都不遞送。 亦即,交易會完整地確定或是完整地捨棄。 您也可以在事務中包含偏移量,以建立讀取、處理訊息並將訊息寫入Kafka應用程式。
程式碼 Snippet
這些程式碼片段在較高層次上說明了所涉及的概念。 如需完整範例,請參閱 GitHub中的 Event Streams 範例。
若要將消費者連接至 Event Streams,您必須建立服務認證。 如需相關資訊,請參閱連接 Event Streams。
在生產者程式碼中,您首先需要建置一組配置內容。 與Event Streams的所有連線均透過使用 TLS 以及使用者和密碼驗證進行保護,因此您至少需要這些屬性。 將 BOOTSTRAP_ENDPOINTS、USER 和 PASSWORD 替換為您自己的服務憑證中的內容:
Properties props = new Properties();
props.put("bootstrap.servers", BOOTSTRAP_ENDPOINTS);
props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"USER\" password=\"PASSWORD\";");
props.put("security.protocol", "SASL_SSL");
props.put("sasl.mechanism", "PLAIN");
props.put("ssl.protocol", "TLSv1.2");
props.put("ssl.enabled.protocols", "TLSv1.2");
props.put("ssl.endpoint.identification.algorithm", "HTTPS");
若要傳送訊息,您還需要為鍵和值指定序列化器,例如:
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
這些序列化程式必須符合消費者使用的解除序列化程式。
然後,使用KafkaProducer發送訊息,其中每條訊息都由ProducerRecord表示。 完成後,別忘了關閉 KafkaProducer。 此程式碼只會傳送訊息,但不等待以查看傳送是否成功。 訊息會傳送至主題 T1,索引鍵為字串 key,值為字串 value。
Producer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
producer.close();
send() 方法非同步,會傳回一個可用來確認其完成的 Future:
Future<RecordMetadata> f = producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
// Do some other stuff
// Now wait for the result of the send
RecordMetadata rm = f.get();
long offset = rm.offset;
或者,您可以在發送訊息時提供回調:
producer.send(new ProducerRecord<String,String>("T1","key","value", new Callback() {
public void onCompletion(RecordMetadata metadata, Exception exception) {
// This is called when the send completes, either successfully or with an exception
}
});
如需相關資訊,請參閱 Kafka 用戶端的 Javadoc。{: external}。