เรียกใช้ไปป์ไลน์
ทำให้กราฟเป็นรูปธรรมและเรียกใช้
เรียกใช้ไปป์ไลน์ เป็นบทเรียน Scala for Backend Engineering & Functional Programming ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Scala for Backend Engineering & Functional Programming และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Scala for Backend Engineering & Functional Programming มีบทเรียนทั้งหมด 4 บทเรียน
จากพิมพ์เขียวสู่การทำงาน
จนถึงตอนนี้ไปป์ไลน์ยังเป็นเพียงพิมพ์เขียว การ สร้างค่าจริง คือกระบวนการเปลี่ยนพิมพ์เขียวนั้นให้เป็นแอกเตอร์ที่ทำงานและเคลื่อนย้ายข้อมูลจริง
จะยังไม่มีอะไรเกิดขึ้นจนกว่าคุณจะเรียกใช้กราฟอย่างชัดเจน ซึ่งทำให้ Akka Streams สามารถประกอบและนำกลับมาใช้ใหม่ได้
ActorSystem
การสร้างค่าจริงต้องใช้ ActorSystem ซึ่งจัดหาเธรดและตัวจัดส่งที่รองรับสเตจของสตรีม ใน Akka รุ่นใหม่ ระบบนี้ยังทำหน้าที่เป็นตัวสร้างค่าจริงโดยปริยายด้วย
โดยทั่วไป ActorSystem หนึ่งรายการจะให้บริการทั้งแอปพลิเคชันและสตรีมที่ทำงานพร้อมกันจำนวนมาก
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
วิธีที่ตรงที่สุดในการรัน Source คือ runWith ซึ่งเชื่อมต่อ Sink และสร้างค่าจริงในขั้นตอนเดียว พร้อมคืนค่าที่สร้างขึ้นเมื่อรันของ Sink นั้น
ในที่นี้ผลลัพธ์คือ Future[Int] ซึ่งจะเสร็จสมบูรณ์พร้อมผลรวมเมื่อสตรีมทำงานเสร็จ
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))run บน RunnableGraph
หากคุณสร้าง RunnableGraph ที่ปิดสมบูรณ์แล้วด้วย to หรือ toMat ให้เรียก run() เพื่อสร้างค่าจริง ค่าที่คืนมาคือค่าที่สร้างขึ้นเมื่อรันซึ่งกราฟเก็บไว้
วิธีนี้แยกการสร้างไปป์ไลน์ออกจากการทำงานจริงได้อย่างชัดเจน
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val graph =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)
val result: Future[Int] = graph.run()ตัวดำเนินการ run แบบอำนวยความสะดวก
Sources มีทางลัด ได้แก่ runForeach, runFold และ runReduce ซึ่งแต่ละรายการจะเชื่อมต่อ Sink ที่เกี่ยวข้องและรันทันที
เหมาะสำหรับการดำเนินการปลายทางทั่วไปบน Source โดยเขียนโค้ดได้กระชับ
import scala.concurrent.Future
val printed: Future[akka.Done] =
Source(1 to 10).runForeach(println)
val sum: Future[Int] =
Source(1 to 10).runFold(0)(_ + _)การทำงานกับ Future ของผลลัพธ์
Sink ปลายทางจะคืนค่าเป็น Future ซึ่งเสร็จสมบูรณ์เมื่อสตรีมสิ้นสุดหรือล้มเหลว ลงทะเบียนการเรียกกลับด้วย onComplete เพื่อจัดการเมื่อสำเร็จหรือเกิดข้อผิดพลาด
ใช้ตัวจัดส่งของ ActorSystem เป็น ExecutionContext โดยปริยายสำหรับการเรียกกลับเหล่านี้
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}ไปป์ไลน์ที่ใกล้เคียงการใช้งานจริง
ไปป์ไลน์ข้อมูลทั่วไปจะอ่านข้อมูลจากแหล่งข้อมูล แปลงข้อมูลด้วยโฟลว์ ดำเนินการอินพุต/เอาต์พุตแบบอะซิงโครนัสด้วย mapAsync แบ่งข้อมูลเป็นชุดด้วย grouped และเขียนข้อมูลลงปลายทาง
แต่ละขั้นตอนมีขนาดเล็ก และทั้งหมดจะถูกทำให้ทำงานจริงด้วย run เพียงครั้งเดียว
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)การเริ่มสตรีมที่ล้มเหลวใหม่
เพื่อเพิ่มความทนทานต่อความขัดข้อง ให้ครอบแหล่งข้อมูลหรือโฟลว์ด้วย RestartSource.withBackoff เพื่อให้ความขัดข้องชั่วคราว เช่น การเชื่อมต่อหลุด กระตุ้นการเริ่มใหม่โดยอัตโนมัติพร้อมการหน่วงเวลาแบบทวีคูณ
วิธีนี้ช่วยให้ไปป์ไลน์นำเข้าข้อมูลที่ทำงานเป็นเวลานานยังทำงานต่อไปได้โดยไม่ต้องควบคุมด้วยตนเอง
import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._
val resilient = RestartSource.withBackoff(
RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)การปิดสตรีมอย่างราบรื่นด้วย KillSwitch
KillSwitch ช่วยให้โค้ดภายนอกหยุดสตรีมที่กำลังทำงานได้อย่างเรียบร้อย ให้แทรก KillSwitches.single ผ่าน viaMat และเก็บค่าที่ได้จากการทำให้ทำงานจริงไว้ เพื่อเรียก shutdown() ในภายหลัง
สิ่งนี้จำเป็นสำหรับสตรีมที่ทำงานเป็นเวลานานและต้องหยุดเมื่อแอปพลิเคชันปิดตัวลง
import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep
val (switch, done) =
source
.viaMat(KillSwitches.single)(Keep.right)
.toMat(Sink.ignore)(Keep.both)
.run()
// later: switch.shutdown()การคืนทรัพยากร
เมื่อแอปพลิเคชันปิดตัวลง ให้เรียก terminate กับ ActorSystem เพื่อคืนเธรดของระบบ เชื่อมการเรียก terminate ไว้หลังฟิวเจอร์ที่ระบุว่าสตรีมทำงานเสร็จสิ้น เพื่อให้การปิดระบบเป็นไปอย่างเป็นระเบียบ
การปล่อย ActorSystem ค้างไว้จะทำให้ JVM ยังคงทำงานและทำให้ทรัพยากรยังคงเปิดใช้งานอยู่
done.onComplete { _ =>
system.terminate()
}การใช้ตัวสร้างการทำงานซ้ำ
การทำให้พิมพ์เขียวเดียวกันทำงานจริงหลายครั้งจะสร้างสตรีมที่กำลังทำงานแยกจากกัน โดยใช้ทรัพยากรของ ActorSystem ร่วมกัน ส่วนพิมพ์เขียวเองยังคงไม่เปลี่ยนแปลงและไม่มีผลข้างเคียง
วิธีนี้ทำให้สามารถกำหนดไปป์ไลน์ไว้เพียงครั้งเดียว แล้วเรียกใช้ด้วย run ตามต้องการสำหรับงานใหม่แต่ละงานที่เข้ามา
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionตรวจสอบความเข้าใจ
พิจารณาว่าอะไรบ้างที่จำเป็นเพื่อให้สตรีมประมวลผลองค์ประกอบต่าง ๆ ได้จริง
สรุปทบทวน
การเรียกใช้ไปป์ไลน์หมายถึงการทำให้พิมพ์เขียวทำงานจริงด้วย ActorSystem ผ่าน run, runWith หรือตัวดำเนินการอำนวยความสะดวก โดยแต่ละวิธีจะคืนผลลัพธ์เป็น Future
คุณได้เรียนรู้ไปป์ไลน์หลายขั้นตอนที่ใกล้เคียงการใช้งานจริง การเริ่มใหม่โดยอัตโนมัติพร้อมการหน่วงเวลา การปิดระบบอย่างราบรื่นผ่าน KillSwitch การล้างทรัพยากรด้วย system.terminate() และการนำพิมพ์เขียวที่ไม่เปลี่ยนแปลงกลับมาใช้ซ้ำอย่างปลอดภัยในการเรียกใช้งานแต่ละครั้งที่เป็นอิสระต่อกัน
เรียนรู้ Scala ด้วย AI tutor — ฟรี
เขียนและเรียกใช้โค้ดจริงในเบราว์เซอร์ของคุณ รับความช่วยเหลือทันทีจาก AI tutor 24/7 และเรียนรู้ต่อจากที่คุณหยุดบนเว็บหรือในแอป
- คอร์ส
- 39
- บทเรียน
- 143
คำถามที่พบบ่อย
บทเรียน “เรียกใช้ไปป์ไลน์” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “เรียกใช้ไปป์ไลน์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ 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 ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน
บทเรียน “เรียกใช้ไปป์ไลน์” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน Scala for Backend Engineering & Functional Programming นี้ได้ไหม
ได้ บทเรียน Scala for Backend Engineering & Functional Programming ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- แหล่งข้อมูล โฟลว์ และซิงก์
- แปลงสตรีม
- แรงดันย้อนกลับ
- เรียกใช้ไปป์ไลน์