Menjalankan Pipeline
Wujudkan dan jalankan graf.
Menjalankan Pipeline adalah pelajaran Scala for Backend Engineering & Functional Programming gratis di CoddyKit. Ini adalah pelajaran 4 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.
Dari Cetak Biru ke Eksekusi
Sampai sejauh ini, pipeline masih berupa cetak biru murni. Materialisasi adalah proses yang mengubah cetak biru tersebut menjadi aktor yang berjalan dan benar-benar memindahkan data.
Tidak ada yang terjadi sampai Anda menjalankan graf secara eksplisit. Inilah yang membuat Akka Streams dapat disusun dan digunakan kembali.
ActorSystem
Materialisasi memerlukan ActorSystem, yang menyediakan thread dan dispatcher untuk menjalankan tahap-tahap stream. Dalam Akka modern, sistem ini juga berperan sebagai materializer implisit.
Satu ActorSystem biasanya melayani seluruh aplikasi dan banyak stream yang berjalan bersamaan.
import akka.actor.ActorSystem
implicit val system: ActorSystem =
ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContextrunWith
Cara paling langsung untuk menjalankan Source adalah runWith, yang memasang Sink dan mematerialisasinya dalam satu langkah, lalu mengembalikan nilai materialisasi Sink tersebut.
Di sini, hasilnya adalah Future[Int] yang selesai dengan jumlah total ketika stream berakhir.
import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future
val total: Future[Int] =
Source(1 to 100).runWith(Sink.fold(0)(_ + _))run pada RunnableGraph
Jika Anda sudah membuat RunnableGraph tertutup dengan to atau toMat, panggil run() untuk mematerialisasinya. Nilai yang dikembalikan adalah nilai materialisasi apa pun yang dipertahankan graf.
Hal ini memisahkan pembuatan pipeline dari eksekusi secara rapi.
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()Operator run Praktis
Source menyediakan pintasan: runForeach, runFold, dan runReduce masing-masing memasang Sink yang sesuai lalu langsung menjalankannya.
Operator tersebut ringkas untuk operasi terminal umum pada 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)(_ + _)Bekerja dengan Future Hasil
Sink terminal mengembalikan Future yang selesai ketika stream berakhir atau gagal. Daftarkan callback dengan onComplete untuk menanggapi keberhasilan atau kesalahan.
Gunakan dispatcher milik ActorSystem sebagai ExecutionContext implisit untuk callback tersebut.
import scala.util.{Success, Failure}
total.onComplete {
case Success(value) => println(s"Sum = $value")
case Failure(ex) => println(s"Failed: ${ex.getMessage}")
}Pipeline yang Realistis
Pipeline data yang umum membaca dari Sumber, melakukan transformasi dengan Alur, menjalankan I/O asinkron dengan mapAsync, mengelompokkan dalam batch dengan grouped, lalu menulis ke Tujuan.
Setiap tahap berukuran kecil, dan seluruhnya dimaterialisasi dengan satu run.
val done =
lineSource
.map(parse)
.mapAsync(4)(validate)
.grouped(500)
.runWith(bulkWriteSink)Memulai Ulang Aliran yang Gagal
Untuk meningkatkan ketahanan, bungkus Sumber atau Alur dengan RestartSource.withBackoff agar kegagalan sementara (seperti koneksi yang terputus) memicu mulai ulang otomatis dengan penundaan eksponensial.
Dengan demikian, pipeline penyerapan data yang berjalan lama tetap aktif tanpa pengawasan manual.
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)Mematikan dengan Anggun menggunakan KillSwitch
KillSwitch memungkinkan kode eksternal menghentikan aliran yang sedang berjalan dengan bersih. Sisipkan KillSwitches.single melalui viaMat, lalu simpan nilai hasil materialisasinya untuk memanggil shutdown() nanti.
Hal ini penting untuk aliran berumur panjang yang harus berhenti saat aplikasi dimatikan.
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()Melepaskan Sumber Daya
Saat aplikasi keluar, hentikan ActorSystem untuk membebaskan thread-nya. Rangkai pemanggilan terminate setelah Future penyelesaian aliran agar proses pematian berlangsung tertib.
ActorSystem yang tidak dilepaskan akan membuat JVM tetap aktif dan mempertahankan sumber daya dalam keadaan terbuka.
done.onComplete { _ =>
system.terminate()
}Menggunakan Kembali Materializer
Mematerialisasi cetak biru yang sama beberapa kali akan membuat aliran independen yang berjalan dan berbagi sumber daya ActorSystem. Cetak biru itu sendiri tetap tidak dapat diubah dan bebas efek samping.
Dengan demikian, Anda dapat mendefinisikan sebuah pipeline sekali, lalu menjalankannya sesuai kebutuhan untuk setiap pekerjaan yang masuk.
val blueprint =
Source(1 to 5).toMat(Sink.seq)(Keep.right)
val run1 = blueprint.run()
val run2 = blueprint.run() // independent executionPemeriksaan Singkat
Pertimbangkan hal-hal yang diperlukan agar suatu aliran benar-benar memproses elemen.
Rangkuman
Menjalankan pipeline berarti mematerialisasi cetak biru dengan sebuah ActorSystem melalui run, runWith, atau operator praktis, yang masing-masing mengembalikan hasil Future.
Anda telah mempelajari pipeline realistis dengan beberapa tahap, mulai ulang otomatis dengan penundaan, pematian dengan anggun melalui KillSwitch, pembersihan sumber daya dengan system.terminate(), serta penggunaan kembali cetak biru yang tidak dapat diubah secara aman dalam beberapa proses independen.
Pertanyaan yang Sering Diajukan
Apakah pelajaran “Menjalankan Pipeline” gratis?
Ya — teks lengkap “Menjalankan Pipeline” 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 “Menjalankan Pipeline”?
Wujudkan dan jalankan graf. 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 4 dari 4.
Berapa lama pelajaran “Menjalankan Pipeline” 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
- Source, Flow, dan Sink
- Mengubah Aliran
- Tekanan Balik
- Menjalankan Pipeline