0Pricing
Scala for Backend Engineering & Functional Programming · Pelajaran

Mengubah Aliran

Petakan dan saring data yang mengalir.

Mengubah Aliran adalah pelajaran Scala for Backend Engineering & Functional Programming gratis di CoddyKit. Ini adalah pelajaran 2 dari 4. Kamu bisa membaca pelajaran lengkapnya di bawah secara gratis — lalu praktikkan langsung di browser dengan editor kode bawaan dan tutor AI 24/7. Ini adalah bagian dari jalur belajar Scala for Backend Engineering & Functional Programming, dan progresmu tersinkronisasi di web dan aplikasi CoddyKit. Kursus Scala for Backend Engineering & Functional Programming mencakup 4 pelajaran total.

Operator sebagai Transformasi

Akka Streams menyediakan sekumpulan operator yang kaya pada Source dan Flow, yang menyerupai API koleksi Scala tetapi berjalan secara asinkron dan mematuhi backpressure.

Setiap operator mengembalikan cetak biru baru, sehingga transformasi dapat disusun secara deklaratif sebelum stream benar-benar dijalankan.

map dan filter

map menerapkan fungsi sinkron ke setiap elemen; filter membuang elemen yang tidak memenuhi predicate. Keduanya merupakan andalan transformasi per elemen.

Keduanya mempertahankan urutan serta meneruskan penyelesaian dan kegagalan ke hilir.

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

mapConcat untuk Satu-ke-Banyak

Jika satu input harus menghasilkan beberapa output, gunakan mapConcat. Operator ini menerima fungsi yang mengembalikan iterable lalu meratakan hasilnya ke dalam stream.

Mengembalikan koleksi kosong secara efektif akan membuang elemen tersebut.

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 dan sliding

grouped(n) mengelompokkan elemen berurutan menjadi Seq yang berisi paling banyak n item, yang berguna untuk penulisan massal ke basis data. sliding(n) menghasilkan jendela yang saling tumpang tindih.

Pengelompokan mengurangi beban per elemen dalam pipeline yang banyak melakukan I/O.

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

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

scan dan fold

scan menghasilkan akumulator yang sedang berjalan setelah setiap elemen, sehingga membentuk stream keadaan yang terus berubah. fold hanya menghasilkan nilai akumulasi akhir sekali setelah upstream selesai.

Gunakan scan untuk penghitung langsung dan fold untuk agregat terminal.

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

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

mapAsync untuk Pekerjaan Asinkron

mapAsync(parallelism) memanggil fungsi yang mengembalikan Future dan menghasilkan hasilnya secara berurutan, dengan menjalankan hingga parallelism Future secara bersamaan.

Gunakan ini untuk pemanggilan asinkron seperti pencarian basis data atau permintaan HTTP ketika urutan penting.

import scala.concurrent.Future

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

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

mapAsyncUnordered

mapAsyncUnordered berperilaku seperti mapAsync, tetapi menghasilkan setiap hasil segera setelah selesai dan mengabaikan urutan input.

Operator ini dapat meningkatkan throughput ketika tahap hilir tidak memedulikan urutan, karena Future yang lambat tidak lagi menghambat Future yang lebih cepat.

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

Transformasi Berkeadaan dengan statefulMapConcat

Untuk transformasi per elemen yang memerlukan keadaan lokal yang dapat diubah, statefulMapConcat membuat keadaan baru untuk setiap materialisasi dan mengembalikan iterable berisi output.

Ini adalah cara aman untuk menyimpan penghitung atau buffer tanpa berbagi keadaan antarjalannya stream.

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

Operator Berbasis Waktu

Stream dapat melakukan transformasi berdasarkan waktu maupun isi. throttle membatasi laju emisi, groupedWithin mengelompokkan berdasarkan ukuran atau waktu yang berlalu, dan takeWithin membatasi durasi.

Operator-operator ini penting untuk membatasi laju API eksternal.

import scala.concurrent.duration._

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

Menangani Kesalahan dalam Transformasi

Pengecualian yang dilemparkan di dalam operator secara bawaan akan menggagalkan seluruh stream. Strategi supervisi dapat memilih resume (membuang elemen yang bermasalah) atau restart tahap tersebut.

Pasang strategi itu dengan withAttributes pada Flow.

import akka.stream.{ActorAttributes, Supervision}

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

Menyusun Flow

Flow kecil dapat disusun menjadi Flow yang lebih besar dengan via, sehingga menghasilkan satu Flow yang dapat digunakan kembali. Hal ini membuat setiap transformasi tetap terfokus dan dapat diuji secara mandiri.

Flow hasil susunan memiliki tipe input dari tahap pertama dan tipe output dari tahap terakhir.

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

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

Pemeriksaan Singkat

Pertimbangkan transformasi asinkron dan jaminan urutannya.

Ringkasan

Anda telah menjelajahi operator transformasi: map/filter per elemen, mapConcat satu-ke-banyak, pengelompokan dengan grouped, akumulasi dengan scan dan fold, serta pekerjaan asinkron melalui mapAsync.

Anda juga mempelajari transformasi berkeadaan, operator berbasis waktu seperti throttle, strategi supervisi untuk menangani kesalahan, dan cara Flow disusun dengan via. Selanjutnya: bagaimana backpressure menjaga tahap-tahap ini tetap aman.

Pertanyaan yang Sering Diajukan

Apakah pelajaran “Mengubah Aliran” gratis?

Ya — teks lengkap “Mengubah Aliran” gratis dibaca di sini di web. Untuk praktiknya secara interaktif (editor kode bawaan dan tutor AI 24/7) dan buka sisa kursus Scala for Backend Engineering & Functional Programming, upgrade ke CoddyKit PRO. Kursus Scala for Backend Engineering & Functional Programming mencakup 4 pelajaran total.

Apa yang akan aku pelajari di “Mengubah Aliran”?

Petakan dan saring data yang mengalir. Kamu berlatih Scala for Backend Engineering & Functional Programming dengan kode praktik yang langsung kamu jalankan di browser, dan tutor AI 24/7 menjawab pertanyaanmu saat kamu mengerjakan pelajaran ini.

Apakah aku perlu pengalaman untuk memulai Scala for Backend Engineering & Functional Programming?

Tidak diperlukan pengalaman sebelumnya. Scala for Backend Engineering & Functional Programming di CoddyKit dirancang untuk pemula hingga pelajar tingkat lanjut, jadi kamu bisa memulai di sini atau dari awal dan belajar sesuai kecepatan kamu sendiri. Ini adalah pelajaran 2 dari 4.

Berapa lama pelajaran “Mengubah Aliran” memakan waktu?

Sebagian besar pelajaran CoddyKit memakan waktu sekitar 5–10 menit. Setiap pelajaran ringkas dan interaktif, jadi kamu membuat kemajuan stabil dan melanjutkan dari tempat kamu tinggalkan di web dan aplikasi.

Bisakah aku menulis dan menjalankan kode dalam pelajaran Scala for Backend Engineering & Functional Programming ini?

Ya. Setiap pelajaran Scala for Backend Engineering & Functional Programming menyertakan editor kode bawaan, jadi kamu menulis dan menjalankan kode nyata langsung di browser dan mendapatkan umpan balik AI instan — tidak diperlukan penyiapan lokal.

Semua pelajaran dalam kursus ini

  1. Source, Flow, dan Sink
  2. Mengubah Aliran
  3. Tekanan Balik
  4. Menjalankan Pipeline
← Kembali ke Scala for Backend Engineering & Functional Programming