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

แปลงสตรีม

แมปและกรองข้อมูลที่ไหลผ่าน

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

ตัวดำเนินการในฐานะการแปลงข้อมูล

Akka Streams มีตัวดำเนินการสำหรับ Source และ Flow ให้ใช้มากมาย โดยมีรูปแบบคล้าย API ของคอลเลกชันใน Scala แต่ทำงานแบบอะซิงโครนัสและเคารพกลไกควบคุมแรงดันย้อนกลับ

ตัวดำเนินการแต่ละตัวจะคืนพิมพ์เขียวใหม่ ทำให้สามารถประกอบการแปลงข้อมูลแบบประกาศได้ก่อนที่สตรีมจะเริ่มทำงาน

map และ filter

map ใช้ฟังก์ชันแบบทำงานพร้อมกันกับสมาชิกทุกตัว ส่วน filter จะทิ้งสมาชิกที่ไม่ผ่านเพรดิเคต ทั้งสองเป็นเครื่องมือหลักสำหรับการแปลงข้อมูลทีละสมาชิก

ทั้งสองรักษาลำดับข้อมูล และส่งต่อสถานะเสร็จสิ้นกับข้อผิดพลาดไปยังปลายทาง

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

mapConcat สำหรับหนึ่งเป็นหลาย

เมื่ออินพุตหนึ่งรายการควรสร้างเอาต์พุตหลายรายการ ให้ใช้ mapConcat โดยจะรับฟังก์ชันที่คืนค่าคอลเลกชันแบบวนซ้ำได้ แล้วคลี่ผลลัพธ์รวมเข้าไปในสตรีม

การคืนค่าคอลเลกชันว่างจะเทียบเท่ากับการทิ้งสมาชิกนั้น

val explode: Flow[String, String, akka.NotUsed] =
  Flow[String].mapConcat(line => line.split(",").toList)

val words = Source(List("a,b", "c,d,e"))
  .via(explode)

grouped และ sliding

grouped(n) จัดสมาชิกที่อยู่ติดกันเป็นชุดในรูป Seq โดยมีได้สูงสุด n รายการ เหมาะสำหรับการเขียนข้อมูลจำนวนมากลงฐานข้อมูล ส่วน sliding(n) จะส่งหน้าต่างข้อมูลที่มีส่วนทับซ้อนกัน

การจัดชุดช่วยลดค่าใช้จ่ายต่อสมาชิกในไปป์ไลน์ที่ทำงานกับระบบรับส่งข้อมูลอย่างหนัก

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan และ fold

scan จะส่งค่าตัวสะสมที่กำลังเปลี่ยนแปลงหลังสมาชิกแต่ละตัว จึงได้สตรีมสถานะที่อัปเดตต่อเนื่อง ส่วน fold จะส่งเฉพาะค่าที่สะสมสุดท้ายเพียงครั้งเดียวเมื่อข้อมูลต้นทางทำงานเสร็จ

ใช้ scan สำหรับตัวนับแบบเรียลไทม์ และใช้ fold สำหรับผลรวมปลายทาง

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

mapAsync สำหรับงานอะซิงโครนัส

mapAsync(parallelism) เรียกฟังก์ชันที่คืนค่าเป็น Future และส่งผลลัพธ์ตาม ลำดับเดิม โดยทำงานพร้อมกันได้สูงสุด parallelism Future

ใช้สำหรับการเรียกแบบอะซิงโครนัส เช่น การค้นหาข้อมูลในฐานข้อมูลหรือคำขอ HTTP ที่ต้องรักษาลำดับ

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered ทำงานคล้าย mapAsync แต่จะส่งผลลัพธ์ทันทีที่แต่ละรายการเสร็จ โดยไม่สนใจลำดับของอินพุต

วิธีนี้อาจเพิ่มอัตราการประมวลผลได้เมื่อสเตจถัดไปไม่สนใจลำดับ เพราะ Future ที่ช้าจะไม่ขัดขวางรายการที่เสร็จเร็วกว่าอีกต่อไป

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

การแปลงที่มีสถานะด้วย statefulMapConcat

สำหรับการแปลงทีละสมาชิกที่ต้องใช้สถานะภายในซึ่งเปลี่ยนแปลงได้ statefulMapConcat จะสร้างสถานะใหม่ทุกครั้งที่สร้างค่าจริง และคืนค่ารายการเอาต์พุตที่วนซ้ำได้

นี่เป็นวิธีที่ปลอดภัยในการเก็บตัวนับหรือบัฟเฟอร์โดยไม่ใช้สถานะร่วมกันระหว่างการรันสตรีม

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

ตัวดำเนินการตามเวลา

สตรีมสามารถแปลงข้อมูลตามเวลาได้เช่นเดียวกับตามเนื้อหา throttle จำกัดอัตราการส่งข้อมูล, groupedWithin จัดชุดตามขนาดหรือเวลาที่ผ่านไป และ takeWithin จำกัดระยะเวลา

สิ่งเหล่านี้จำเป็นอย่างยิ่งสำหรับการจำกัดอัตราการเรียก API ภายนอก

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

การจัดการข้อผิดพลาดในการแปลงข้อมูล

ข้อยกเว้นที่ถูกโยนภายในตัวดำเนินการจะทำให้สตรีมทั้งหมดล้มเหลวตามค่าเริ่มต้น กลยุทธ์กำกับดูแลสามารถเลือกให้ resume (ทิ้งสมาชิกที่มีปัญหา) หรือ restart สเตจแทนได้

แนบกลยุทธ์ดังกล่าวด้วย withAttributes บน Flow

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

การประกอบ Flows

สามารถประกอบ Flows ขนาดเล็กให้เป็น Flow ขนาดใหญ่ด้วย via ซึ่งจะได้ Flow เดียวที่นำกลับมาใช้ใหม่ได้ วิธีนี้ทำให้การแปลงแต่ละส่วนมีหน้าที่ชัดเจนและทดสอบแยกกันได้

Flow ที่ประกอบแล้วจะมีชนิดอินพุตเป็นชนิดของส่วนแรก และชนิดเอาต์พุตเป็นชนิดของส่วนสุดท้าย

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

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

พิจารณาการแปลงแบบอะซิงโครนัสและการรับประกันเรื่องลำดับ

สรุปทบทวน

คุณได้สำรวจตัวดำเนินการแปลงข้อมูล ได้แก่ map/filter สำหรับการประมวลผลทีละสมาชิก, mapConcat สำหรับหนึ่งเป็นหลาย, grouped สำหรับการจัดชุด, scan และ fold สำหรับการสะสม และงานอะซิงโครนัสผ่าน mapAsync

คุณยังได้เห็นการแปลงที่มีสถานะ ตัวดำเนินการตามเวลาอย่าง throttle กลยุทธ์กำกับดูแลสำหรับข้อผิดพลาด และวิธีประกอบ Flows ด้วย via ต่อไป: กลไกควบคุมแรงดันย้อนกลับช่วยให้สเตจเหล่านี้ปลอดภัยได้อย่างไร

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

บทเรียน “แปลงสตรีม” ฟรีหรือไม่

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