Scala for Backend Engineering & Functional Programming · บทเรียน

เรียกใช้ไปป์ไลน์

ทำให้กราฟเป็นรูปธรรมและเรียกใช้

บทเรียน 4 จาก 413 ขั้นตอน

เรียกใช้ไปป์ไลน์ เป็นบทเรียน 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 // ExecutionContext

runWith

วิธีที่ตรงที่สุดในการรัน 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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

บทเรียนทั้งหมดในหลักสูตรนี้

  1. แหล่งข้อมูล โฟลว์ และซิงก์
  2. แปลงสตรีม
  3. แรงดันย้อนกลับ
  4. เรียกใช้ไปป์ไลน์
← กลับไปที่ Scala for Backend Engineering & Functional Programming