แปลงสตรีม
แมปและกรองข้อมูลที่ไหลผ่าน
แปลงสตรีม เป็นบทเรียน 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)(_ + _) // 15mapAsync สำหรับงานอะซิงโครนัส
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 ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ