Che cos'è PySpark?

PySpark è l'API Python ufficiale per Apache Spark, un framework di computing distribuito open source progettato per il trattamento dati su larga scala e il machine learning. Il vantaggio principale di PySpark è che consente a data engineer e data scientist di scrivere codice Python familiare che distribuisce automaticamente i carichi di lavoro su molti cluster di computer. Si tratta di un cambiamento significativo rispetto all'esecuzione di script Python a nodo singolo, che possono facilmente bloccarsi a causa di limitazioni di memoria.

Questa architettura distribuita è progettata per gestire set di dati di grandi dimensioni, eseguire trasformazioni di dati complesse in memoria e gestire pipeline ETL (Extract, Transform, Load) di livello aziendale. In questo modo, PySpark fornisce una solida base per l'elaborazione dei big data.

Perché usare PySpark per l'elaborazione di big data?

PySpark rivoluziona il data engineering passando dalla scalabilità verticale (acquisto di macchine più grandi e potenti) alla scalabilità orizzontale (distribuzione del carico di lavoro su più macchine). Questo approccio è particolarmente efficace per la gestione di grandi volumi di dati non strutturati. Ecco alcuni dei vantaggi principali:

  • Calcolo in memoria: PySpark può memorizzare nella cache i dati in memoria su diversi nodi, a differenza dei sistemi più vecchi come Hadoop MapReduce. Ciò accelera notevolmente gli algoritmi iterativi e le attività di machine learning che devono accedere più volte agli stessi dati. 
  • Valutazione lazy: PySpark utilizza un approccio "lazy", il che significa che non esegue immediatamente le trasformazioni. Invece, attende che venga chiamata un'azione. Ciò consente al Catalyst Optimizer integrato di determinare il modo più efficiente per eseguire le attività senza richiedere un intervento manuale. 
  • Tolleranza di errore: PySpark si basa su una struttura di dati chiamata Resilient Distributed Datasets (RDD). Gli RDD tengono traccia di come vengono trasformati i dati, quindi se un nodo worker non funziona durante un'attività, il sistema può recuperare automaticamente i dati persi.

I moduli principali di PySpark

PySpark è più di un singolo strumento, è un ecosistema completo di moduli che consentono agli sviluppatori di gestire tutto, dalla pulizia dei dati di base al machine learning avanzato e allo streaming in tempo reale, il tutto all'interno di un unico framework.

PySpark Core e RDD

PySpark Core è la base dell'intero sistema. Fornisce le funzionalità di base, inclusa l'API di basso livello per i set di dati distribuiti resilienti (RDD). Sebbene gli RDD offrano un controllo dettagliato e siano a tolleranza di errore, la maggior parte delle applicazioni moderne utilizza astrazioni di livello superiore come i DataFrame, che offrono una migliore ottimizzazione.

Spark SQL e DataFrame

Il "dataframe pyspark" è il modulo con cui lavora la maggior parte degli sviluppatori. Un DataFrame PySpark è una raccolta distribuita di dati organizzati in colonne con nome, simile a una tabella in un database. Questa struttura consente a Catalyst Optimizer di creare piani di esecuzione altamente efficienti che possono superare le prestazioni del codice Python standard.

Machine learning con MLlib

PySpark MLlib è una libreria di machine learning scalabile che fornisce API di alto livello per creare, addestrare ed eseguire il deployment di pipeline di machine learning. Ciò include attività come la regressione, il clustering e la classificazione, tutte eseguite su set di dati distribuiti senza la necessità di spostare i dati su altri sistemi.

Structured Streaming

Questo modulo distingue tra elaborazione batch e analisi in tempo reale. Structured Streaming è un motore di elaborazione dei flussi di dati a tolleranza di errore che consente agli sviluppatori di eseguire query SQL continue e incrementali su flussi di dati in tempo reale provenienti da origini come Apache Kafka o socket TCP. Puoi anche creare pipeline di analisi dei flussi di dati utilizzando esattamente la stessa API DataFrame che utilizzi per i dati statici, il che mantiene il tuo codebase semplice e facile da gestire.

Comprendere il passaggio da Pandas a PySpark

Gli strumenti tradizionali di elaborazione dei dati, come la libreria standard Pandas, operano su una singola macchina, caricando interi set di dati nella memoria (RAM) di quella macchina. Questo approccio funziona bene per i set di dati più piccoli, ma presenta rapidamente problemi quando si ha a che fare con i set di dati di grandi dimensioni comuni negli ambienti aziendali. Questi sistemi a nodo singolo non possono essere scalati orizzontalmente, il che porta a errori di esecuzione e di memoria.

Sebbene PySpark richieda una configurazione iniziale più complessa, risolve direttamente questo collo di bottiglia della memoria. Offre una via di mezzo che combina la sintassi facile da imparare di Python con la potenza dell'elaborazione distribuita di Spark. PySpark partiziona automaticamente i dati, li distribuisce su più nodi worker ed esegue le attività in parallelo.

Operazioni di base e ottimizzazione dei dati

Scrivere codice PySpark efficiente significa comprendere come i dati vengono trasformati ed elaborati in un cluster. Ecco alcuni dei concetti chiave per lavorare con i DataFrame:

  • Trasformazioni e azioni: è importante comprendere la differenza tra le trasformazioni (come map(), filter() e join()), che creano nuovi DataFrame in modo lazy, e le azioni (come count(), show() e collect()), che attivano il calcolo effettivo sul cluster.
  • Pulizia dei dati e gestione dei valori nulli: i data engineer utilizzano PySpark per gestire i valori mancanti, correggere i tipi di dati e normalizzare i set di dati disordinati su larga scala prima che i dati vengano utilizzati per l'analisi.
  • Partizionamento e gestione della memoria: per ottimizzare le prestazioni, gli sviluppatori possono gestire il modo in cui i dati vengono partizionati e possono trasmettere tabelle più piccole a tutti i nodi per evitare costosi data shuffling durante le operazioni di unione.

Scrittura di codice PySpark per pipeline di dati aziendali

L'API DataFrame di PySpark consente agli sviluppatori di passare dalla scrittura di codice procedurale a nodo singolo a codice dichiarativo distribuito. Sebbene la sintassi sia simile alle librerie Python come Pandas, l'esecuzione è fondamentalmente diversa. Il codice PySpark crea un piano logico che Catalyst Optimizer valuta ed esegue in parallelo nel cluster. Ecco i pattern principali per la creazione di una pipeline di dati:

  • Acquisizione dei dati e creazione di DataFrame: si inizia inizializzando una SparkSession e poi leggendo grandi set di dati da origini come l'archiviazione di oggetti o i database relazionali in un DataFrame distribuito utilizzando comandi come spark.read.format().load(). 
  • Trasformazioni e aggregazioni dichiarative: i tecnici possono concatenare i metodi per manipolare gli schemi di dati. Ciò include la selezione di colonne, il filtraggio di righe e l'esecuzione di aggregazioni distribuite con groupBy() e agg() per elaborare milioni di record senza causare problemi di memoria.
  • Esecuzione di Spark SQL nativo: PySpark consente una grande interoperabilità permettendo agli sviluppatori di registrare un DataFrame come vista temporanea con createOrReplaceTempView. Da lì, possono eseguire query SQL ANSI standard direttamente all'interno dei loro script Python utilizzando spark.sql().
  • Generazione di codice assistita dall'AI: i moderni IDE basati su cloud possono accelerare lo sviluppo di PySpark. Ad esempio, gli assistenti di programmazione AI sono disponibili nei notebook di BigQuery Studio e possono generare automaticamente operazioni complesse di PySpark DataFrame da prompt in linguaggio naturale.

Domande frequenti

Ecco alcune domande comuni su PySpark:

Pandas viene eseguito su una singola macchina ed elabora i dati in un singolo spazio di memoria, il che lo rende ideale per set di dati di piccole dimensioni. PySpark, d'altra parte, è un motore di computing distribuito che partiziona i dati su un cluster di macchine, il che lo rende necessario quando i set di dati sono troppo grandi per essere contenuti nella RAM di un singolo computer.

Gli RDD sono strutture di dati di basso livello che non hanno uno schema definito, il che li rende flessibili per dati complessi e non strutturati, ma più lenti da elaborare. I DataFrame sono basati su RDD, ma applicano uno schema (righe e colonne) che consente a Catalyst Optimizer di Spark di migliorare automaticamente le prestazioni delle query per i dati strutturati.

La valutazione pigra è la strategia di PySpark che consiste nel registrare le istruzioni logiche (trasformazioni) senza calcolare immediatamente i dati. Attende un comando di "azione" prima di calcolare il piano di esecuzione più efficiente ed eseguire il calcolo.

Sì, PySpark è ampiamente utilizzato per creare ed eseguire processi ETL (Extract, Transform, Load). La sua capacità di connettersi a diverse origini dati, eseguire trasformazioni potenti su set di dati di grandi dimensioni e caricare i dati in vari sistemi lo rende una scelta eccellente per l'ETL. Tuttavia, è più di un semplice strumento ETL: è un framework completo che include anche librerie per il machine learning (MLlib), l'elaborazione dei flussi e l'analisi dei grafi, il che lo rende una piattaforma per uso generico per i big data.

I vantaggi di PySpark nell'era dell'AI generativa

Man mano che le aziende si orientano verso l'AI generativa e i modelli linguistici di grandi dimensioni (LLM), la sfida più grande è spesso la preparazione dei dati, non l'addestramento del modello. PySpark fornisce la potenza di elaborazione distribuita necessaria per importare, pulire e trasformare petabyte di dati non strutturati nel contesto strutturato di alta qualità di cui i moderni modelli di AI hanno bisogno per funzionare correttamente.

Feature engineering su vasta scala

I DataFrame PySpark consentono ai machine learning engineer di eseguire l'estrazione e la vettorizzazione di caratteristiche complesse su miliardi di righe in parallelo, riducendo significativamente il tempo necessario per preparare i set di dati per l'addestramento o il perfezionamento dei modelli.


Elaborazione di dati non strutturati per RAG

Le pipeline di Retrieval-Augmented Generation (RAG) e gli agenti AI si basano su grandi quantità di dati di testo e log. L'architettura distribuita di PySpark è perfettamente adatta per l'analisi, la suddivisione in blocchi e la pulizia di questi dati non strutturati prima che vengano incorporati nei database vettoriali.


Integrazione perfetta con l'ecosistema AI

PySpark funge da collegamento vitale tra i data lake non elaborati e i framework di AI avanzati. Consente ai team di inserire dati puliti e distribuiti direttamente in librerie di deep learning come PyTorch o TensorFlow, o in piattaforme di AI aziendali come Gemini, senza dover esportare i dati in altri sistemi.


Casi d'uso di PySpark

PySpark sta rapidamente diventando lo standard per l'architettura delle pipeline di dati in molti settori principali. Consentendo il passaggio dall'elaborazione dei dati locale al computing distribuito, PySpark aiuta le aziende ad affrontare sfide complesse in termini di infrastruttura e analisi.

Rilevamento delle frodi in tempo reale

Gli istituti finanziari utilizzano PySpark Streaming per proteggere le transazioni. Le applicazioni di streaming possono analizzare continuamente i log delle transazioni in tempo reale e confrontarli con i modelli di rischio storici creati con MLlib per segnalare attività fraudolente prima che si traducano in una perdita finanziaria.

Motori per suggerimenti per l'e-commerce

Le aziende di vendita al dettaglio utilizzano PySpark per creare esperienze utente dinamiche. Un workflow tipico prevede l'importazione di petabyte di dati di clickstream degli utenti, la loro pulizia con operazioni DataFrame e l'addestramento di modelli di filtro collaborativo per personalizzare i prezzi e i consigli sui prodotti in tempo reale.

Manutenzione predittiva e telemetria IoT

Le fabbriche intelligenti utilizzano PySpark per migliorare la produzione industriale. Le pipeline di dati possono estrarre dati in tempo reale dai sensori su macchinari pesanti, elaborare questi dati non strutturati su larga scala e prevedere guasti hardware per programmare la manutenzione in modo proattivo.

Pipeline di dati ETL su larga scala

Gli amministratori di database e gli architetti di sistema utilizzano PySpark per risolvere i colli di bottiglia dei database. Ad esempio, potrebbero sostituire i vecchi script batch con PySpark per estrarre enormi set di dati da varie località di archiviazione cloud, trasformare i dati in memoria e quindi caricare i dati puliti e ottimizzati in un data warehouse centrale.

Scalare i carichi di lavoro PySpark su Google Cloud

Google Cloud offre un ambiente potente per la scalabilità dei workload PySpark dai prototipi locali alla produzione aziendale completa. Managed Service for Apache Spark fornisce un hub centrale in cui gli sviluppatori possono eseguire PySpark senza dover configurare o ottimizzare manualmente i cluster. La piattaforma è ottimizzata sia per i job batch a lunga esecuzione che per lo streaming in tempo reale, offrendo sia cluster gestiti che modalità di deployment serverless.

La piattaforma fornisce anche strumenti e integrazioni specializzati per gli sviluppatori. I team di progettazione possono collegare senza problemi le loro pipeline PySpark con l'ecosistema di dati più ampio di Google Cloud, incluse le query dirette a BigQuery per l'analisi su scala petabyte. I team possono anche connettere le loro pipeline di dati PySpark alla Gemini Enterprise Agent Platform per creare, gestire ed eseguire il deployment di modelli ML e agenti AI avanzati basati su dati aziendali puliti e distribuiti.

Guida alla creazione e al deployment di agenti autonomi sul tuo lakehouse

Un approccio moderno al deployment di agenti AI autonomi prevede l'esecuzione diretta su una data lakehouse unificata. Il framework di deployment Antigravity consente agli architetti di eseguire il deployment di questi agenti dove già risiedono i dati, anziché spostare set di dati massicci elaborati con PySpark in ambienti LLM esterni. Questa architettura riduce al minimo la latenza di rete e i costi in uscita, mantenendo al contempo una rigorosa governance dei dati.

Per colmare il divario tra l'elaborazione distribuita dei dati e l'AI agentica, i team di engineering possono utilizzare Data Agent Kit per creare flussi di lavoro multi-agente in grado di eseguire query in modo nativo su DataFrame PySpark e tabelle lakehouse. Il Model Context Protocol (MCP) funge da livello sicuro che consente a questi agenti di recuperare dinamicamente il contesto da strumenti e API aziendali esterni senza compromettere la sicurezza dell'ambiente lakehouse.

Fai il prossimo passo

Inizia a creare su Google Cloud con 300 $ di crediti gratuiti e oltre 20 prodotti Always Free.

Google Cloud