पाठ 7 / 25
Batching और Compression
Throughput के लिए linger.ms, batch.size और compression ट्यून करें।
थोड़ी latency देकर बहुत throughput पाएँ
Producers हर record अलग नहीं भेजते। वे प्रति partition records को batches में इकट्ठा करते हैं। linger.ms बताता है कि batch भरने के लिए कितना इंतज़ार करना है (0 तुरंत भेजता है; 5 से 20 ms अक्सर throughput काफ़ी सुधारता है), और batch.size bytes में अधिकतम batch आकार है। compression.type (lz4, zstd, snappy, gzip) पूरे batches को compress करता है, जिससे network और disk उपयोग घटता है; lz4 और zstd अच्छी गति व अनुपात के लिए लोकप्रिय हैं। Compression producer पर होता है और disk पर भी वैसा ही रहता है, और consumers decompress करते हैं। अपने डेटा से मापें, क्योंकि नतीजे record के आकार और सामग्री पर निर्भर हैं।
Python में producer (उदाहरण)
यह confluent-kafka package (pip install confluent-kafka) उपयोग करता है। इसे यहाँ चलाया नहीं गया; config keys मानक नाम हैं। flush() बाक़ी भेजे जाने का इंतज़ार करता है।
from confluent_kafka import Producer
p = Producer({
"bootstrap.servers": "localhost:9092",
"acks": "all",
"enable.idempotence": True,
"linger.ms": 10,
"compression.type": "lz4",
})
def on_delivery(err, msg):
if err: print("failed:", err)
else: print(f"ok {msg.topic()}[{msg.partition()}]@{msg.offset()}")
for i in range(5):
p.produce("orders", key=f"user-{i}", value=f"order {i}", on_delivery=on_delivery)
p.flush()Delivery नतीजे हमेशा जाँचें
produce() asynchronous है। Delivery callback अनदेखा करें या बाहर निकलने से पहले flush() न बुलाएँ, तो आपके कोड में कोई error आए बिना records खो सकते हैं।
त्वरित जाँच: 10 जैसा छोटा `linger.ms` आम तौर पर क्या करता है?
- Batches को थोड़ा भरने देता है, थोड़ी latency के बदले throughput बढ़ाता है
- Retries बंद करता है
- डेटा encrypt करता है
- पुराने records हटाता है
Answer
Batches को थोड़ा भरने देता है, थोड़ी latency के बदले throughput बढ़ाता है — थोड़ा रुकने से ज़्यादा records batch में जुड़ते हैं, जिससे दक्षता बढ़ती है।