Scala for Backend Engineering & Functional Programming · Lezione

Eseguire una pipeline

Materializzi ed esegua un grafo.

Lezione 4 di 413 passaggi

Eseguire una pipeline è una lezione Scala for Backend Engineering & Functional Programming gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento Scala for Backend Engineering & Functional Programming, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso Scala for Backend Engineering & Functional Programming include 4 lezioni in totale.

Dal blueprint all'esecuzione

Finora la pipeline è stata un blueprint puro. La materializzazione è il processo che trasforma quel blueprint in actor in esecuzione, capaci di trasferire effettivamente i dati.

Non accade nulla finché non si esegue esplicitamente il grafo: è questo che rende Akka Streams componibile e riutilizzabile.

L'ActorSystem

La materializzazione richiede un ActorSystem, che fornisce i thread e il dispatcher su cui si basano gli stage dello stream. Nelle versioni moderne di Akka, il sistema funge anche da materializer implicito.

In genere un ActorSystem serve un'intera applicazione e numerosi stream concorrenti.

import akka.actor.ActorSystem

implicit val system: ActorSystem =
  ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContext

runWith

Il modo più diretto per eseguire un Source è runWith, che collega un Sink e materializza il grafo in un unico passaggio, restituendo il valore materializzato di quel Sink.

In questo caso il risultato è un Future[Int] che si completa con la somma al termine dello stream.

import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future

val total: Future[Int] =
  Source(1 to 100).runWith(Sink.fold(0)(_ + _))

Eseguire un RunnableGraph

Se ha già costruito un RunnableGraph chiuso con to o toMat, chiami run() per materializzarlo. Il valore restituito è il valore materializzato conservato dal grafo.

In questo modo la costruzione della pipeline viene separata con chiarezza dalla sua esecuzione.

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()

Operatori di esecuzione rapida

I Source offrono alcune scorciatoie: runForeach, runFold e runReduce collegano ciascuno il Sink corrispondente ed eseguono immediatamente lo stream.

Sono soluzioni concise per le operazioni terminali più comuni su un 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)(_ + _)

Gestire il Future del risultato

I sink terminali restituiscono un Future che si completa quando lo stream termina o fallisce. Registri le callback con onComplete per reagire al successo o all'errore.

Per queste callback usi il dispatcher dell'ActorSystem come ExecutionContext implicito.

import scala.util.{Success, Failure}

total.onComplete {
  case Success(value) => println(s"Sum = $value")
  case Failure(ex)    => println(s"Failed: ${ex.getMessage}")
}

Una pipeline realistica

Una pipeline di dati tipica legge da una Source, trasforma i dati con i Flow, esegue operazioni di I/O asincrone con mapAsync, raggruppa gli elementi in batch con grouped e scrive in un Sink.

Ogni fase è piccola e l'intera pipeline viene materializzata con una sola chiamata a run.

val done =
  lineSource
    .map(parse)
    .mapAsync(4)(validate)
    .grouped(500)
    .runWith(bulkWriteSink)

Riavvio degli stream non riusciti

Per garantire la resilienza, racchiuda una Source o un Flow in RestartSource.withBackoff, in modo che gli errori temporanei, come una connessione interrotta, attivino un riavvio automatico con backoff esponenziale.

In questo modo le pipeline di acquisizione dati di lunga durata rimangono attive senza supervisione manuale.

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)

Arresto ordinato con KillSwitch

Un KillSwitch consente al codice esterno di arrestare in modo ordinato uno stream in esecuzione. Inserisca KillSwitches.single tramite viaMat e conservi il relativo valore materializzato per poter chiamare in seguito shutdown().

Questo è essenziale per gli stream di lunga durata che devono arrestarsi quando l'applicazione viene terminata.

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()

Rilascio delle risorse

Quando l'applicazione termina, termini l'ActorSystem per liberare i relativi thread. Esegua la chiamata a terminate dopo il Future di completamento dello stream, così l'arresto avverrà in modo ordinato.

Lasciare un ActorSystem non terminato mantiene attiva la JVM e lascia aperte le risorse.

done.onComplete { _ =>
  system.terminate()
}

Riutilizzo del Materializer

Materializzare più volte lo stesso blueprint crea stream indipendenti in esecuzione che condividono le risorse dell'ActorSystem. Il blueprint rimane immutabile e privo di effetti collaterali.

Questo consente di definire una pipeline una sola volta ed eseguirla su richiesta per ogni nuovo job.

val blueprint =
  Source(1 to 5).toMat(Sink.seq)(Keep.right)

val run1 = blueprint.run()
val run2 = blueprint.run() // independent execution

Verifica rapida

Consideri ciò che è necessario affinché uno stream elabori effettivamente gli elementi.

Riepilogo

Eseguire una pipeline significa materializzare un blueprint con un ActorSystem tramite run, runWith o gli operatori di utilità, ciascuno dei quali restituisce un Future.

Ha visto pipeline realistiche composte da più fasi, riavvii automatici con backoff, arresti ordinati tramite KillSwitch, pulizia delle risorse con system.terminate() e il riutilizzo sicuro di un blueprint immutabile in esecuzioni indipendenti.

Gratis per iniziare

Impara Scala con un tutor IA — gratis

Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.

Corsi
39
Lezioni
143

Domande Frequenti

La lezione «Eseguire una pipeline» è gratuita?

Sì — il testo completo di «Eseguire una pipeline» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso Scala for Backend Engineering & Functional Programming, passa a CoddyKit PRO. Il corso Scala for Backend Engineering & Functional Programming include 4 lezioni in totale.

Cosa imparerò in «Eseguire una pipeline»?

Materializzi ed esegua un grafo. Eserciti Scala for Backend Engineering & Functional Programming con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare Scala for Backend Engineering & Functional Programming?

Non è richiesta alcuna esperienza precedente. Scala for Backend Engineering & Functional Programming su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «Eseguire una pipeline»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione Scala for Backend Engineering & Functional Programming?

Sì. Ogni lezione Scala for Backend Engineering & Functional Programming include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Source, Flow e Sink
  2. Trasformare gli stream
  3. Backpressure
  4. Eseguire una pipeline
← Torna a Scala for Backend Engineering & Functional Programming