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

แหล่งข้อมูล โฟลว์ และซิงก์

องค์ประกอบพื้นฐานของการสตรีม

แหล่งข้อมูล โฟลว์ และซิงก์ เป็นบทเรียน Scala for Backend Engineering & Functional Programming ฟรีบน CoddyKit นี่คือบทเรียนที่ 1 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Scala for Backend Engineering & Functional Programming และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Scala for Backend Engineering & Functional Programming มีบทเรียนทั้งหมด 4 บทเรียน

องค์ประกอบพื้นฐานสามอย่าง

Akka Streams จำลอง pipeline ข้อมูลเป็นกราฟของขั้นตอนการประมวลผล ขั้นตอนเชิงเส้นหลักสามอย่างคือ Source (สร้างองค์ประกอบ), Flow (แปลงองค์ประกอบ) และ Sink (รับองค์ประกอบไปใช้งาน)

Source มีเอาต์พุตหนึ่งรายการ Sink มีอินพุตหนึ่งรายการ และ Flow มีอินพุตหนึ่งรายการกับเอาต์พุตหนึ่งรายการพอดี การเชื่อมต่อสิ่งเหล่านี้เข้าด้วยกันจะอธิบายว่า จะเกิดอะไรขึ้น ไม่ใช่ เกิดขึ้นเมื่อใด

การกำหนด Source

Source[Out, Mat] ส่งองค์ประกอบชนิด Out และเปิดให้ใช้ค่าที่สร้างขึ้นจริงชนิด Mat Source ที่ง่ายที่สุดมาจากคอลเลกชันหรือช่วงค่าที่อยู่ในหน่วยความจำ

จนกว่าจะเรียกใช้สตรีม Source เป็นเพียงต้นแบบที่เปลี่ยนแปลงไม่ได้ และสามารถนำกลับมาใช้ใหม่ได้อย่างอิสระ

import akka.stream.scaladsl.Source

val numbers: Source[Int, akka.NotUsed] =
  Source(1 to 100)

val single: Source[String, akka.NotUsed] =
  Source.single("hello")

การกำหนด Sink

Sink[In, Mat] รับองค์ประกอบชนิด In ไปใช้งาน ค่าที่สร้างขึ้นจริงมักเก็บผลลัพธ์ของการรับข้อมูล เช่น Future ที่เสร็จสมบูรณ์เมื่อสตรีมทำงานเสร็จ

Sink.foreach จะเรียกใช้ผลข้างเคียงกับแต่ละองค์ประกอบ ส่วน Sink.fold จะสะสมเป็นผลลัพธ์เดียว

import akka.stream.scaladsl.Sink
import scala.concurrent.Future

val printSink: Sink[Int, Future[akka.Done]] =
  Sink.foreach(println)

val sumSink: Sink[Int, Future[Int]] =
  Sink.fold(0)(_ + _)

การกำหนด Flow

Flow[In, Out, Mat] อยู่ระหว่าง Source กับ Sink และแปลงแต่ละองค์ประกอบ Flow สามารถนำกลับมาใช้ได้ด้วยตัวเอง และสามารถประกอบกันก่อนเชื่อมต่อกับปลายทางใด ๆ

ในที่นี้ Flow จะคูณจำนวนเต็มด้วยสองและแปลงเป็นสตริง

import akka.stream.scaladsl.Flow

val doubleToString: Flow[Int, String, akka.NotUsed] =
  Flow[Int]
    .map(_ * 2)
    .map(n => s"value=$n")

การเชื่อมต่อ Source กับ Sink

ตัวดำเนินการ via ใช้เชื่อม Flow เข้ากับ Source และ to ใช้เชื่อม Sink การเชื่อมต่อ Source เข้ากับ Sink โดยตรงด้วย to จะสร้าง RunnableGraph ซึ่งเป็นต้นแบบที่ปิดสมบูรณ์และเรียกใช้งานได้

ยังไม่มีองค์ประกอบใดเคลื่อนที่ การทำงานนี้เพียงอธิบายโครงสร้างการเชื่อมต่อ

import akka.stream.scaladsl.{Source, Sink, RunnableGraph}

val graph: RunnableGraph[akka.NotUsed] =
  Source(1 to 10).to(Sink.foreach(println))

via: การแทรก Flow

ใช้ via เพื่อแทรกโฟลว์ลงในไปป์ไลน์ โดย Source.via(flow) จะให้ Source ใหม่ที่ชนิดข้อมูลเอาต์พุตตรงกับเอาต์พุตของ Flow

การเรียก via ต่อกันช่วยให้คุณสร้างไปป์ไลน์การแปลงข้อมูลขนาดยาวจากชิ้นส่วน Flow ขนาดเล็กที่ทดสอบได้ง่าย

val pipeline =
  Source(1 to 10)
    .via(Flow[Int].filter(_ % 2 == 0))
    .via(Flow[Int].map(_ * 10))
    .to(Sink.foreach(println))

ความปลอดภัยของชนิดข้อมูลตลอดทุกสเตจ

คอมไพเลอร์จะตรวจสอบว่าชนิดข้อมูลเอาต์พุตของแต่ละสเตจตรงกับชนิดข้อมูลอินพุตของสเตจถัดไป โดย Source[Int] ไม่สามารถเชื่อมต่อกับ Sink[String] ได้ หากไม่มี Flow ระหว่างกลางที่แปลงชนิดข้อมูล

การตรวจสอบแบบสแตติกนี้จะตรวจพบข้อผิดพลาดในการเชื่อมต่อไปป์ไลน์ก่อนเริ่มทำงานจริง

// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
  Source(1 to 3)
    .via(Flow[Int].map(_.toString))
    .to(Sink.foreach[String](println))

ค่าที่สร้างขึ้นเมื่อรัน

พิมพ์เขียวแต่ละรายการจะมี ค่าที่สร้างขึ้นเมื่อรัน: ตัวจัดการที่สร้างขึ้นเมื่อสตรีมทำงาน Sink.fold จะสร้าง Future ของผลลัพธ์ โดยค่าเริ่มต้น การรวมสเตจจะเก็บค่าที่สร้างขึ้นเมื่อรันจากฝั่งซ้ายสุด (NotUsed สำหรับแหล่งข้อมูลทั่วไป)

ใช้ toMat และ Keep เพื่อเลือกค่าจากฝั่งที่ต้องการ

import akka.stream.scaladsl.Keep
import scala.concurrent.Future

val g: RunnableGraph[Future[Int]] =
  Source(1 to 100)
    .toMat(Sink.fold(0)(_ + _))(Keep.right)

องค์ประกอบที่นำกลับมาใช้ใหม่ได้

เนื่องจาก Sources, Flows และ Sinks เป็นค่าที่เปลี่ยนแปลงไม่ได้ คุณจึงกำหนดค่าเหล่านี้ครั้งเดียวแล้วนำกลับมาใช้ในหลายไปป์ไลน์ได้ วิธีนี้ส่งเสริมให้สร้างคลังสเตจประมวลผลขนาดเล็กที่มีชื่อเรียกชัดเจน

สามารถนำ Flow ที่กำหนดไว้สำหรับการแยกวิเคราะห์ไปใส่ในทั้งไปป์ไลน์ไฟล์และไปป์ไลน์ HTTP ได้

val parse: Flow[String, Int, akka.NotUsed] =
  Flow[String].map(_.trim.toInt)

val fromFile  = lines.via(parse)
val fromHttp  = requestBody.via(parse)

ตัวสร้าง Source ที่ใช้บ่อย

Akka Streams มีโรงงานสร้าง Source หลายแบบ ได้แก่ Source.single, Source.repeat, Source.tick สำหรับการส่งข้อมูลตามเวลา, Source.future จาก Future และ Source.empty

การเลือกตัวสร้างที่เหมาะสมช่วยสื่อเจตนาของผู้ผลิตข้อมูลได้อย่างชัดเจน

import scala.concurrent.duration._

val ticks  = Source.tick(0.seconds, 1.second, "tick")
val onceF  = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")

ตัวสร้าง Sink ที่ใช้บ่อย

ในทำนองเดียวกัน Sink มีทั้ง Sink.head (สมาชิกแรกในรูป Future), Sink.seq (รวบรวมทั้งหมดเป็น Seq), Sink.ignore (อ่านจนหมดแล้วทิ้ง) และ Sink.last

สำหรับไปป์ไลน์ที่ส่งผลลัพธ์กลับไปยังโค้ดของคุณ Sink.seq และ Sink.fold เป็นตัวเลือกที่ใช้บ่อยที่สุด

import scala.concurrent.Future

val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]]    = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignore

ตรวจสอบความเข้าใจอย่างรวดเร็ว

ทดสอบความเข้าใจของคุณเกี่ยวกับชนิดของสเตจแบบเชิงเส้น

สรุปทบทวน

คุณได้เรียนรู้หน่วยประกอบแบบเชิงเส้นสามชนิด ได้แก่ Source ทำหน้าที่ผลิตข้อมูล, Flow แปลงข้อมูล และ Sink รับข้อมูล ทั้งหมดเป็นพิมพ์เขียวที่เปลี่ยนแปลงไม่ได้และนำกลับมาใช้ใหม่ได้ โดยเชื่อมต่อกันด้วย via และ to

การเชื่อมต่อ Source เข้ากับ Sink จะได้ RunnableGraph ที่มีค่าซึ่งสร้างขึ้นเมื่อรัน แต่จะยังไม่เคลื่อนย้ายข้อมูลจนกว่าจะเรียกใช้ ต่อไปคุณจะได้เรียนรู้การแปลงสตรีมด้วยตัวดำเนินการที่มีความสามารถมากขึ้น

คำถามที่พบบ่อย

บทเรียน “แหล่งข้อมูล โฟลว์ และซิงก์” ฟรีหรือไม่

ใช่ — ข้อความเต็มของ “แหล่งข้อมูล โฟลว์ และซิงก์” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ 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 ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 1 จากทั้งหมด 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