แรงดันย้อนกลับ
จัดการผู้ผลิตที่ทำงานเร็วอย่างปลอดภัย
แรงดันย้อนกลับ เป็นบทเรียน Scala for Backend Engineering & Functional Programming ฟรีบน CoddyKit นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Scala for Backend Engineering & Functional Programming และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Scala for Backend Engineering & Functional Programming มีบทเรียนทั้งหมด 4 บทเรียน
แรงดันย้อนกลับคืออะไร
แรงดันย้อนกลับเป็นกลไกควบคุมการไหลที่ป้องกันไม่ให้ผู้ผลิตข้อมูลที่เร็วทำให้ผู้บริโภคที่ช้ารับภาระเกินกำลัง แทนที่จะบัฟเฟอร์ข้อมูลอย่างไม่จำกัดหรือทิ้งข้อมูล ผู้บริโภคจะส่งสัญญาณว่าตนรองรับได้มากเพียงใด
Akka Streams ใช้มาตรฐาน Reactive Streams ซึ่งความต้องการจะไหลย้อนขึ้นต้นทาง และสมาชิกข้อมูลจะไหลไปยังปลายทาง
การไหลที่ขับเคลื่อนด้วยความต้องการ
แต่ละสเตจจะส่งข้อมูลก็ต่อเมื่อสเตจถัดไปส่งสัญญาณ ความต้องการ แล้ว Sink จะร้องขอสมาชิกจำนวน N รายการ และความต้องการนั้นจะแพร่ย้อนกลับไปจนถึง Source ซึ่งจะผลิตข้อมูลตรงตามที่ร้องขอ
โพรโทคอลแบบดึงนี้ทำให้ผู้ผลิตไม่ส่งข้อมูลมากกว่าที่ผู้บริโภคจะประมวลผลได้
เหตุผลที่เรื่องนี้สำคัญต่อไปป์ไลน์
หากไม่มีแรงดันย้อนกลับ ผู้บริโภค Kafka ที่เร็วซึ่งส่งข้อมูลไปยังฐานข้อมูลที่ช้าจะสะสมระเบียนที่กำลังประมวลผลอยู่หลายล้านรายการ จนหน่วยความจำเต็มและกระบวนการทำงานล้มเหลว
แรงดันย้อนกลับจะลดความเร็วของต้นทางให้เท่ากับสเตจที่ช้าที่สุดโดยอัตโนมัติ ทำให้การใช้หน่วยความจำคงที่แม้มีภาระงานสูง
// Fast source, slow sink: backpressure slows the source
val g =
Source(1 to 1000000)
.map(_ * 2)
.to(slowDatabaseSink)การบัฟเฟอร์ภายใน
ระหว่างขอบเขตอะซิงโครนัส Akka Streams จะเก็บบัฟเฟอร์ภายในขนาดเล็ก (ค่าเริ่มต้น 16 รายการ) เพื่อรองรับการส่งข้อมูลเป็นชุดสั้น ๆ ทำให้สเตจไม่จำเป็นต้องทำงานสอดคล้องกันทุกสมาชิก
เมื่อบัฟเฟอร์เต็ม แรงดันย้อนกลับจะเริ่มทำงาน และต้นทางจะหยุดผลิตข้อมูลจนกว่าจะมีพื้นที่ว่าง
import akka.stream.Attributes
val buffered =
Flow[Int]
.map(identity)
.addAttributes(Attributes.inputBuffer(initial = 32, max = 32))บัฟเฟอร์แบบชัดเจนด้วยกลยุทธ์ล้น
ตัวดำเนินการ buffer จะแทรกบัฟเฟอร์แบบชัดเจนตามขนาดที่เลือก พร้อม OverflowStrategy ที่กำหนดว่าจะทำอย่างไรเมื่อบัฟเฟอร์เต็ม
วิธีนี้ทำให้คุณแลกการใช้หน่วยความจำกับความสามารถในการแยกความเร็วของผู้ผลิตและผู้บริโภคออกจากกัน
import akka.stream.OverflowStrategy
val withBuffer =
Source(1 to 1000)
.buffer(size = 100, OverflowStrategy.backpressure)กลยุทธ์เมื่อบัฟเฟอร์ล้น
กลยุทธ์ต่าง ๆ ได้แก่ backpressure (ลดความเร็วต้นทาง), dropHead/dropTail (ทิ้งข้อมูลเก่าที่สุดหรือใหม่ที่สุด), dropBuffer, dropNew และ fail (ยุติการทำงานพร้อมข้อผิดพลาด)
กลยุทธ์การทิ้งข้อมูลเหมาะกับข้อมูลสด เช่น ค่าจากเซนเซอร์ ซึ่งสามารถทิ้งค่าที่ล้าสมัยได้อย่างปลอดภัย
import akka.stream.OverflowStrategy
val latestWins =
liveTicks.buffer(1, OverflowStrategy.dropHead)
val strict =
liveTicks.buffer(50, OverflowStrategy.fail)ใช้ conflate เพื่อสรุปข้อมูล
เมื่อผู้บริโภคทำงานช้า conflate จะรวมสมาชิกที่รออยู่เป็นหนึ่งรายการด้วยฟังก์ชันรวม แทนที่จะบัฟเฟอร์สมาชิกทั้งหมด
ตัวอย่างเช่น รวมการอัปเดตตัวเลขหลายรายการให้เป็นผลรวม เพื่อให้ผู้บริโภคเห็นค่ารวมของข้อมูลที่พลาดไปเสมอ
val summarized =
fastMetrics
.conflate((acc, next) => acc + next)
// Slow downstream receives summed batchesใช้ expand เพื่อเติมเต็มความต้องการ
expand เป็นแนวคิดตรงข้ามกับ conflate: เมื่อปลายทางต้องการข้อมูลเร็วกว่าที่ต้นทางผลิตได้ ระบบจะสร้างสมาชิกเพิ่มเติมจากค่าล่าสุดที่เห็น
เหมาะสำหรับการส่งค่าที่อ่านได้ล่าสุดอย่างต่อเนื่องด้วยอัตราคงที่
val repeated =
sensor.expand(last => Iterator.continually(last))
// Downstream always gets the latest sensor valueขอบเขตอะซิงโครนัส
ตามค่าเริ่มต้น สเตจที่หลอมรวมกันจะทำงานบนแอกเตอร์เดียวโดยไม่มีบัฟเฟอร์คั่นกลาง การแทรก async จะวางสเตจไว้บนแอกเตอร์ของตนเอง เพิ่มบัฟเฟอร์และเปิดให้ประมวลผลแบบขนานในไปป์ไลน์
ขอบเขตอะซิงโครนัสคือจุดที่บัฟเฟอร์ของแรงดันย้อนกลับอาศัยอยู่จริง
val pipelined =
Source(1 to 1000)
.map(slowStep).async
.map(anotherSlowStep).async
.to(Sink.ignore)throttle ในฐานะการควบคุมอัตราอย่างชัดเจน
throttle กำหนดอัตราสูงสุดโดยเจตนา และสร้างแรงดันย้อนกลับไปยังต้นทางเพื่อรักษาอัตรานั้น วิธีนี้ช่วยปกป้องบริการภายนอกที่จำกัดอัตรา แม้ว่าผู้บริโภคจะทำงานได้เร็วกว่า
พารามิเตอร์ burst อนุญาตให้มีการพุ่งขึ้นชั่วครู่เหนืออัตราปกติ
import scala.concurrent.duration._
val limited =
requests
.throttle(
elements = 100, per = 1.second, maximumBurst = 20,
akka.stream.ThrottleMode.Shaping)การสังเกตแรงดันย้อนกลับ
คุณตรวจจับแรงดันย้อนกลับได้โดยสังเกตการชะลอตัวของต้นทางหรือวัดการครอบครองบัฟเฟอร์ ตัวดำเนินการ log และแอตทริบิวต์สตรีมของ Akka ช่วยติดตามจุดที่ไปป์ไลน์หยุดชะงัก
บัฟเฟอร์ที่เต็มอยู่ตลอดบ่งชี้สเตจที่ช้าที่สุด ซึ่งเป็นตัวจำกัดอัตราการประมวลผล
val traced =
Source(1 to 100)
.log("after-source")
.map(_ * 2)
.log("after-map")
.to(Sink.ignore)ตรวจสอบความเข้าใจอย่างรวดเร็ว
ลองคิดว่า Akka Streams ป้องกันไม่ให้ผู้ผลิตที่เร็วส่งข้อมูลท่วมผู้บริโภคที่ช้าได้อย่างไร
สรุปทบทวน
แรงดันย้อนกลับเป็นแกนหลักของ Akka Streams ที่ขับเคลื่อนด้วยความต้องการ ผู้บริโภคจะส่งสัญญาณความต้องการย้อนขึ้นต้นทาง เพื่อไม่ให้ผู้ผลิตทำงานเกินกำลัง และช่วยจำกัดการใช้หน่วยความจำ
คุณได้เห็นบัฟเฟอร์ภายใน ตัวดำเนินการ buffer แบบชัดเจนพร้อมกลยุทธ์เมื่อบัฟเฟอร์ล้น การสรุปด้วย conflate การเติมเต็มความต้องการด้วย expand ขอบเขตอะซิงโครนัส และการควบคุมอัตราโดยเจตนาด้วย throttle ต่อไปคุณจะได้รันไปป์ไลน์ที่สมบูรณ์
คำถามที่พบบ่อย
บทเรียน “แรงดันย้อนกลับ” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “แรงดันย้อนกลับ” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส Scala for Backend Engineering & Functional Programming ให้อัปเกรดเป็น CoddyKit PRO คอร์ส Scala for Backend Engineering & Functional Programming มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “แรงดันย้อนกลับ”
จัดการผู้ผลิตที่ทำงานเร็วอย่างปลอดภัย คุณปฏิบัติ Scala for Backend Engineering & Functional Programming ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน Scala for Backend Engineering & Functional Programming หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน Scala for Backend Engineering & Functional Programming บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 3 จากทั้งหมด 4 บทเรียน
บทเรียน “แรงดันย้อนกลับ” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน Scala for Backend Engineering & Functional Programming นี้ได้ไหม
ได้ บทเรียน Scala for Backend Engineering & Functional Programming ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- แหล่งข้อมูล โฟลว์ และซิงก์
- แปลงสตรีม
- แรงดันย้อนกลับ
- เรียกใช้ไปป์ไลน์