我的 Spring Boot 應用程式將每小時從 kafka 代理收聽 100 萬條記錄。每條訊息的整個處理邏輯需要 1-1.5 秒,包括資料庫插入。Broker有64個磁區,這也是我@KafkaListener的并發。
在我每小時收聽大約 50k 條記錄的較低環境中,我當前的代碼只能在一分鐘內處理 90 條記錄。以下是代碼,所有其他配置引數(如 max.poll.records 等)均為默認值:
@KafkaListener(id="xyz-listener", concurrency="64", topics="my-topic")
public void listener(String record) {
// processing logic
}
我確實得到每小時 7-8 次“消費者很可能被踢出小組”。我認為這兩個問題都可以通過隔離偵聽器方法和每條訊息的多執行緒處理來解決,但我不知道該怎么做。
uj5u.com熱心網友回復:
這里有幾點需要考慮。首先,對于單個應用程式而言,64 個消費者似乎有點過多,無法始終如一地處理。
考慮到默認情況下每次輪詢一次獲取500 records每個消費者,如果處理單個批次的時間超過默認的 5 分鐘,您的應用程式可能會過載并導致消費者被踢出組max.poll.timeout.ms。
所以首先,我會考慮scaling the application horizontally讓每個應用程式處理較少數量的磁區/執行緒。
增加吞吐量的第二種方法是使用批處理偵聽器,并按您在此答案中看到的那樣批量處理處理和資料庫插入。
使用這兩者,您應該在每個應用程式中并行處理大量的作業,并且應該能夠實作您想要的吞吐量。
當然,您應該使用不同的數字對每種方法進行負載測驗,以獲得適當的指標。
編輯:解決您的評論,如果您想實作此吞吐量,我還不會放棄批處理。如果您逐行執行資料庫操作,您將需要更多資源才能獲得相同的性能。
如果您的規則引擎不執行任何 I/O,您可以通過它迭代批次中的每條記錄而不會損失性能。
關于資料一致性,可以嘗試一些策略。例如,您可以lock確保即使通過重新平衡,也只有一個實體將在給定時間處理給定批次的記錄 - 或者在 Kafka 中使用重新平衡掛鉤可能有一種更慣用的處理方式。
有了它,您可以在收到記錄時批量加載過濾掉重復/過期記錄所需的所有資訊,通過記憶體中的規則引擎迭代每條記錄,然后批量持久化所有結果,然后釋放鎖。
當然,如果不了解流程的更多細節,就很難提出理想的策略。關鍵是通過這樣做,您應該能夠在每個實體中處理大約 10 倍以上的記錄,所以我肯定會試一試。
轉載請註明出處,本文鏈接:https://www.uj5u.com/houduan/491168.html
