Passa al contenuto principale
Open source

Introduzione alle Arrow UDFs in PySpark: Un sostituto più veloce e leggero per le Pandas UDFs

Definisci UDF più performanti con facilità.

di Ruifeng Zheng e Yicong Huang

  • Introduciamo le UDF Arrow native, che operano direttamente sui dati Arrow, eliminando l'overhead di conversione Pandas/Arrow nelle UDF Pandas per un'esecuzione più rapida e un minore utilizzo di memoria.
  • Descriviamo anche i tipi di UDF Arrow per casi d'uso scalari e di aggregazione, e le UDTF Arrow per trasformazioni table-in, table-out, con esempi di codice sia in Python che in SQL.
  • I benchmark mostrano che le UDF Arrow sono circa il 10% più veloci e utilizzano circa il 40% in meno di memoria rispetto alle UDF Pandas, con un migliore supporto per i tipi di dati complessi.

Introduzione

Le funzioni definite dall'utente (UDF) Python sono un meccanismo di estensibilità essenziale, ma tradizionalmente hanno sofferto di un elevato overhead dovuto all'esecuzione basata su righe. In Apache Spark™, le UDF Pandas hanno affrontato parte di questo problema introducendo la serializzazione basata su Arrow e l'elaborazione batch, migliorando significativamente il throughput rispetto alle UDF Python scalari.

Tuttavia, le UDF Pandas presentano ancora limitazioni fondamentali:

  • La conversione dei dati Pandas/Arrow introduce copie di dati aggiuntive. Gli approcci a copia zero sono possibili solo in alcuni casi specifici. Ad esempio, le colonne con valori NULL attiveranno copie profonde.
  • I tipi di dati complessi non sono ben supportati. Ad esempio, le istanze di StructType nidificate non sono supportate per il tipo di output con casi d'uso di aggregazione.
Flussi di dati dell'esecuzione delle UDF Pandas in Apache Spark

Eliminando la conversione dei dati Pandas/Arrow, le UDF Arrow vengono eseguite più velocemente delle UDF Pandas, consumano meno memoria e offrono un migliore supporto per i tipi di dati.

UDF Arrow Native

Siamo entusiasti di presentare le UDF Arrow Native a partire da Databricks Runtime 18.0 (note di rilascio), un entusiasmante passo avanti per l'esecuzione performante delle UDF.

Le UDF Arrow Native operano direttamente sui dati Arrow senza convertire gli input in oggetti Pandas o NumPy. Ciò preserva il layout colonnare end-to-end, evita copie di dati non necessarie e consente alle UDF di utilizzare l'elaborazione vettorializzata sfruttando il modello di calcolo e memoria nativo di Arrow.

Per definire una UDF Arrow, gli utenti possono utilizzare un nuovo decoratore Python @arrow_udf, con tipo di ritorno specificato e tipo di valutazione opzionale. Ad esempio:

Gli utenti possono anche definirla con il decoratore esistente @udf con suggerimenti di tipo completi. Ad esempio:

Nota: la definizione della funzione dovrebbe includere suggerimenti di tipo per tutti gli argomenti e il valore di ritorno.
Questo design si allinea con le interfacce delle UDF Python scalari, fornendo un'esperienza coerente e intuitiva per gli utenti già familiari con le UDF Python scalari.

Quanto segue dimostra come utilizzare la UDF Arrow:

Utilizzo Python:

Utilizzo SQL:

Forniamo supporto per varianti delle interfacce UDF Arrow. Incluse funzioni scalari, funzioni di aggregazione e funzioni di tabella. Nell'API del dataframe forniamo anche mapInArrow e applyInArrow per utilizzare le UDF Arrow. Le introdurremo una per una.

Funzioni Scalari Arrow

Le funzioni scalari Arrow eseguono trasformazioni riga per riga. Sono l'equivalente Arrow delle UDF Pandas scalari e possono essere utilizzate ovunque sia prevista un'espressione di colonna, come df.select() o df.withColumn(). Sono supportate tre modalità di input: diretta, iteratore e iteratore di array multipli. Le varianti dell'iteratore sono utili quando la UDF richiede un'inizializzazione costosa una tantum (ad esempio, il caricamento di un modello o la compilazione di un pattern regex), poiché il costo di configurazione viene ammortizzato su tutti i batch. In tutti i casi, il numero di righe di output deve corrispondere al numero di righe di input.

  • Da Array ad Array: riceve uno o più pyarrow.Array e restituisce un pyarrow.Array. L'array di input e di output deve avere lo stesso numero di valori.
  • Da Iteratore di Array a Iteratore di Array: riceve un iteratore di pyarrow.Array e restituisce un iteratore di pyarrow.Array. Questo tipo è utile quando l'esecuzione della UDF richiede un'inizializzazione costosa.
  • Da Iteratore di Array Multipli a Iteratore di Array: riceve un iteratore di una tupla di più pyarrow.Array e restituisce un iteratore di pyarrow.Array.

Funzioni di Aggregazione Arrow

Le funzioni di aggregazione Arrow accettano uno o più input pyarrow.Array e restituiscono un valore scalare, riducendo un gruppo di righe in un singolo risultato. Sono l'equivalente Arrow delle UDF Pandas aggregate raggruppate e vengono utilizzate con groupBy().agg() o operazioni Window. Similmente alle funzioni scalari, le funzioni di aggregazione supportano anche tre modalità di input.

Da Array a Scalare: riceve pyarrow.Array e restituisce un valore scalare.

  • Da Iteratore di Array a Scalare: riceve un iteratore di pyarrow.Array e restituisce un valore scalare. Questo è utile per elaborare grandi volumi di dati in operazioni di aggregazione.

Da Iteratore di Array Multipli a Scalare: riceve un iteratore di una tupla di più pyarrow.Array e restituisce un valore scalare. È possibile definire aggregazioni più complesse.

Funzioni di Tabella Arrow

Le funzioni di tabella Arrow, note anche come UDTF Arrow (User-Defined Table Functions), accettano un pyarrow.RecordBatch o più pa.Array come input e producono un pyarrow.Table come output. Questo rappresenta il modello predominante per le trasformazioni table-in, table-out implementate in Python utilizzando l'esecuzione colonnare. Le UDTF Arrow possiedono la capacità di:

  • Restituire più colonne
  • Produrre zero, una o più righe
  • Eseguire trasformazioni di tabella vettorializzate impiegando kernel di calcolo Arrow

Di conseguenza, sono ottimamente adatte per operazioni come il filtraggio, l'espansione delle righe, la ristrutturazione dei dati e la generazione di colonne derivate.

L'interfaccia arrow_udtf è progettata per la semplicità, impiegando una sintassi a decoratore in cui si definisce il tipo di ritorno usando una stringa formattata DDL. In questa configurazione, il metodo eval accetta oggetti PyArrow come input e si prevede che produca tabelle PyArrow o RecordBatches. L'interfaccia supporta due modalità di input. Quando si elaborano argomenti di tabella, al metodo eval viene fornito un oggetto pa.RecordBatch che incapsula tutte le colonne dalla tabella di input:

Per gli argomenti scalari, il metodo riceve oggetti pa.Array, uno per ogni input scalare:

Ecco un altro esempio:

Questa UDTF può funzionare in due modi distinti:

Utilizzo Python:

Utilizzo SQL:

Supporto per DataFrame mapInArrow e applyInArrow

Oltre alle User-Defined Functions (UDF) e alle User-Defined Table Functions (UDTF), PySpark fornisce API di funzione Arrow che facilitano l'applicazione diretta di funzioni native Python ai dati Arrow a livello di DataFrame. Queste API operano in modo analogo alle loro controparti Pandas (mapInPandas, applyInPandas) ma utilizzano pyarrow.RecordBatch e pyarrow.Table invece dei DataFrame Pandas, aggirando così l'overhead di conversione tra i formati Pandas e Arrow.

  • Mappa. DataFrame.mapInArrow trasforma un iteratore di pyarrow.RecordBatch in un altro iteratore di pyarrow.RecordBatch, consentendo operazioni a livello di riga come il filtraggio, la trasformazione o l'espansione.
  • Mappa raggruppata. groupBy().applyInArrow() applica una funzione specificata a ogni gruppo, accettando e restituendo un pyarrow.Table. Questa funzionalità si rivela utile per le trasformazioni per gruppo, come la normalizzazione dei dati.
  • Mappa co-raggruppata. cogroup().applyInArrow() consente il co-raggruppamento di due DataFrame basati su una chiave condivisa, applicando successivamente una funzione a ogni co-gruppo. La funzione riceve due pyarrow.Table input e si prevede che restituisca un singolo pyarrow.Table.

Prestazioni

Eliminando la costosa conversione dei dati Pandas/Arrow, le UDF Arrow generalmente vengono eseguite più velocemente delle UDF Pandas, con un minore utilizzo di memoria. Confrontiamo le due semplici UDF:

L'UDF Arrow è circa il 10% più veloce dell'UDF Pandas, e il profiler di memoria mostra che circa il 40% di memoria viene risparmiato nell'esecuzione.

Conclusione

Databricks Runtime 18.0 introduce le UDF Arrow native, offrendo un'alternativa più veloce e snella alle UDF Pandas per un'esecuzione performante delle UDF Python in PySpark. Operando direttamente sui dati Arrow ed eliminando l'overhead di conversione Pandas/Arrow, le UDF Arrow offrono un'esecuzione circa il 10% più veloce, circa il 40% in meno di utilizzo di memoria e un migliore supporto per i tipi di dati complessi, il tutto con una sintassi a decoratore familiare e intuitiva.

Pronto a esplorare di più? Prova oggi stesso le UDF Arrow native su Databricks come parte di Databricks Runtime 18.0. Per iniziare, sostituisci semplicemente le tue UDF Pandas esistenti con le UDF Arrow. Nella maggior parte dei casi, bastano poche righe di codice per sbloccare guadagni di prestazioni immediati. Consulta la documentazione delle UDF Arrow e la documentazione delle UDTF Arrow per il riferimento API completo e ulteriori esempi.

(Questo post sul blog è stato tradotto utilizzando strumenti basati sull'intelligenza artificiale) Post originale

Ricevi gli ultimi articoli nella tua casella di posta

Iscriviti al nostro blog e ricevi gli ultimi articoli direttamente nella tua casella di posta.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])Da Iteratore di Array Multipli a Iteratore di Array: riceve un iteratore di una tupla di pi\u00f9 pyarrow.Array e restituisce un iteratore di pyarrow.Array.python@arrow_udf(\"long\")\ndef multiply(iterator: Iterator[Tuple[pa.Array, pa.Array]]) -> Iterator[pa.Array]:\n for v1, v2 in iterator:\n yield pa.compute.multiply(v1, v2)Funzioni di Aggregazione ArrowLe funzioni di aggregazione Arrow accettano uno o pi\u00f9 input pyarrow.Array e restituiscono un valore scalare, riducendo un gruppo di righe in un singolo risultato. Sono l'equivalente Arrow delle UDF Pandas aggregate raggruppate e vengono utilizzate con groupBy().agg() o operazioni Window. Similmente alle funzioni scalari, le funzioni di aggregazione supportano anche tre modalit\u00e0 di input. Da Array a Scalare: riceve pyarrow.Array e restituisce un valore scalare. python@arrow_udf(\"struct<m1: double, m2: double>\")\ndef min_max_udf(v: pa.Array) -> pa.Scalar:\n m1 = pa.compute.min(v)\n m2 = pa.compute.max(v)\n t = pa.struct([pa.field(\"m1\", pa.float64()), pa.field(\"m2\", pa.float64())])\n return pa.scalar(value={\"m1\": m1.as_py(), \"m2\": m2.as_py()}, type=t)Da Iteratore di Array a Scalare: riceve un iteratore di pyarrow.Array e restituisce un valore scalare. Questo \u00e8 utile per elaborare grandi volumi di dati in operazioni di aggregazione.python@arrow_udf(\"double\")\ndef streaming_mean(iterator: Iterator[pa.Array]) -> float:\n total_sum = 0.0\n total_count = 0\n for batch in iterator:\n total_sum += pa.compute.sum(batch).as_py()\n total_count += len(batch)\n return total_sum / total_count if total_count > 0 else 0.0Da Iteratore di Array Multipli a Scalare: riceve un iteratore di una tupla di pi\u00f9 pyarrow.Array e restituisce un valore scalare. \u00c8 possibile definire aggregazioni pi\u00f9 complesse.python@arrow_udf(\"double\")\ndef weighted_mean(iterator: Iterator[Tuple[pa.Array, pa.Array]]) -> float:\n weighted_sum = 0.0\n total_weight = 0.0\n for values, weights in iterator:\n weighted_sum += pa.compute.sum(pa.compute.multiply(values, weights)).as_py()\n total_weight += pa.compute.sum(weights).as_py()\n return weighted_sum / total_weight if total_weight > 0 else 0.0Funzioni di Tabella ArrowLe funzioni di tabella Arrow, note anche come UDTF Arrow (User-Defined Table Functions), accettano un pyarrow.RecordBatch o pi\u00f9 pa.Array come input e producono un pyarrow.Table come output. Questo rappresenta il modello predominante per le trasformazioni table-in, table-out implementate in Python utilizzando l'esecuzione colonnare. Le UDTF Arrow possiedono la capacit\u00e0 di:Restituire pi\u00f9 colonneProdurre zero, una o pi\u00f9 righeEseguire trasformazioni di tabella vettorializzate impiegando kernel di calcolo ArrowDi conseguenza, sono ottimamente adatte per operazioni come il filtraggio, l'espansione delle righe, la ristrutturazione dei dati e la generazione di colonne derivate.L'interfaccia arrow_udtf \u00e8 progettata per la semplicit\u00e0, impiegando una sintassi a decoratore in cui si definisce il tipo di ritorno usando una stringa formattata DDL. In questa configurazione, il metodo eval accetta oggetti PyArrow come input e si prevede che produca tabelle PyArrow o RecordBatches. L'interfaccia supporta due modalit\u00e0 di input. Quando si elaborano argomenti di tabella, al metodo eval viene fornito un oggetto pa.RecordBatch che incapsula tutte le colonne dalla tabella di input:python@arrow_udtf(returnType=\"x int, y int\")\nclass ProcessBatch:\n def eval(self, batch: pa.RecordBatch):\n x_array = batch.column('x')\n y_array = batch.column('y')\n result = pa.table({\n 'x': pc.multiply(x_array, 2),\n 'y': pc.add(y_array, 10)\n })\n yield resultPer gli argomenti scalari, il metodo riceve oggetti pa.Array, uno per ogni input scalare:python@arrow_udtf(returnType=\"x int, y int\")\nclass ProcessArrays:\n def eval(self, x: pa.Array, y: pa.Array):\n result = pa.table({\n 'x': x,\n 'y': pc.multiply(y, 2)\n })\n yield resultEcco un altro esempio:python@arrow_udtf(returnType=\"sensor_id string, temp_f double, status string\")\nclass ProcessReadings:\n def eval(self, batch: pa.RecordBatch, threshold: pa.Array):\n min_temp = threshold[0].as_py()\n \n # Vectorized filter: keep readings above threshold\n mask = pc.greater(batch.column(\"temp_c\"), min_temp)\n filtered = pa.table(batch).filter(mask)\n \n # Vectorized transform: Celsius to Fahrenheit\n temp_f = pc.add(pc.multiply(filtered.column(\"temp_c\"), 1.8), 32)\n \n # Vectorized conditional: assign status based on temperature\n status = pc.if_else(\n pc.greater(filtered.column(\"temp_c\"), 35),\n pa.scalar(\"CRITICAL\"),\n pa.scalar(\"WARNING\")\n )\n \n yield pa.table({\n \"sensor_id\": filtered.column(\"sensor_id\"),\n \"temp_f\": temp_f,\n \"status\": status\n })Questa UDTF pu\u00f2 funzionare in due modi distinti:Utilizzo Python:pythonfrom pyspark.sql import functions as F\n\n# Generate sample data\ndf = spark.createDataFrame([\n (\"sensor_1\", 30.0),\n (\"sensor_2\", 36.0),\n (\"sensor_3\", 25.0),\n], [\"sensor_id\", \"temp_c\"])\n\n# Invoke the UDTF within Python\nresult = ProcessReadings(df.asTable(), F.lit(28.0))\n\n\nresult.show()\n# +----------+-------+--------+\n# |sensor_id |temp_f |status |\n# +----------+-------+--------+\n# |sensor_1 |86.0 |WARNING |\n# |sensor_2 |96.8 |CRITICAL|\n# +----------+-------+--------+Utilizzo SQL:python# Register the UDTF for SQL environment access\nspark.udtf.register(\"process_readings\", ProcessReadings)\n\n# Establish a temporary table or view\ndf.createOrReplaceTempView(\"sensor_data\")\n\n# Execute the UDTF in SQL\nspark.sql(\"\"\"\n SELECT *\n FROM process_readings(\n TABLE(SELECT sensor_id, temp_c FROM sensor_data),\n 28.0\n )\n\"\"\").show()\n# +----------+-------+--------+\n# |sensor_id |temp_f |status |\n# +----------+-------+--------+\n# |sensor_1 |86.0 |WARNING |\n# |sensor_2 |96.8 |CRITICAL|\n# +----------+-------+--------+Supporto per DataFrame mapInArrow e applyInArrowOltre alle User-Defined Functions (UDF) e alle User-Defined Table Functions (UDTF), PySpark fornisce API di funzione Arrow che facilitano l'applicazione diretta di funzioni native Python ai dati Arrow a livello di DataFrame. Queste API operano in modo analogo alle loro controparti Pandas (mapInPandas, applyInPandas) ma utilizzano pyarrow.RecordBatch e pyarrow.Table invece dei DataFrame Pandas, aggirando cos\u00ec l'overhead di conversione tra i formati Pandas e Arrow.Mappa. DataFrame.mapInArrow trasforma un iteratore di pyarrow.RecordBatch in un altro iteratore di pyarrow.RecordBatch, consentendo operazioni a livello di riga come il filtraggio, la trasformazione o l'espansione.pythonimport pyarrow as pa\n\ndf = spark.createDataFrame([(1, 21), (2, 30)], (\"id\", \"age\"))\n\ndef filter_func(iterator):\n for batch in iterator:\n yield batch.filter(pa.compute.field(\"id\") == 1)\n\ndf.mapInArrow(filter_func, df.schema).show()\n# +---+---+\n# | id|age|\n# +---+---+\n# | 1| 21|\n# +---+---+Mappa raggruppata. groupBy().applyInArrow() applica una funzione specificata a ogni gruppo, accettando e restituendo un pyarrow.Table. Questa funzionalit\u00e0 si rivela utile per le trasformazioni per gruppo, come la normalizzazione dei dati.pythonimport pyarrow as pa\nimport pyarrow.compute as pc\n\ndf = spark.createDataFrame(\n [(1, 1.0), (1, 2.0), (2, 3.0), (2, 5.0), (2, 10.0)], (\"id\", \"v\"))\n\ndef normalize(table):\n v = table.column(\"v\")\n norm = pc.divide(pc.subtract(v, pc.mean(v)), pc.stddev(v, ddof=1))\n return table.set_column(1, \"v\", norm)\n\ndf.groupby(\"id\").applyInArrow(normalize, schema=\"id long, v double\").show()\n# +---+-------------------+\n# | id| v|\n# +---+-------------------+\n# | 1|-0.7071067811865...|\n# | 1| 0.7071067811865...|\n# | 2|-0.8320502943378...|\n# | 2|-0.2773500981126...|\n# | 2| 1.1094003924504...|\n# +---+-------------------+Mappa co-raggruppata. cogroup().applyInArrow() consente il co-raggruppamento di due DataFrame basati su una chiave condivisa, applicando successivamente una funzione a ogni co-gruppo. La funzione riceve due pyarrow.Table input e si prevede che restituisca un singolo pyarrow.Table.pythonimport pyarrow as pa\n\ndf1 = spark.createDataFrame(\n [(1, 1.0), (2, 2.0), (1, 3.0), (2, 4.0)], (\"id\", \"v1\"))\ndf2 = spark.createDataFrame([(1, \"x\"), (2, \"y\")], (\"id\", \"v2\"))\n\ndef summarize(l, r):\n return pa.Table.from_pydict({\n \"left\": [l.num_rows],\n \"right\": [r.num_rows]\n })\n\ndf1.groupby(\"id\").cogroup(df2.groupby(\"id\")).applyInArrow(\n summarize, schema=\"left long, right long\").show()\n# +----+-----+\n# |left|right|\n# +----+-----+\n# | 2| 1|\n# | 2| 1|\n# +----+-----+PrestazioniEliminando la costosa conversione dei dati Pandas/Arrow, le UDF Arrow generalmente vengono eseguite pi\u00f9 velocemente delle UDF Pandas, con un minore utilizzo di memoria. Confrontiamo le due semplici UDF:python@pandas_udf(\"long\")\ndef multiply_pandas_func(a: pd.Series, b: pd.Series) -> pd.Series:\n return a * b\n\n@arrow_udf(\"long\")\ndef multiply_arrow_func(a: pa.Array, b: pa.Array) -> pa.Array:\n return pa.compute.multiply(a, b)L'UDF Arrow \u00e8 circa il 10% pi\u00f9 veloce dell'UDF Pandas, e il profiler di memoria mostra che circa il 40% di memoria viene risparmiato nell'esecuzione.python============================================================\nProfile of UDF<id=2>\n============================================================\nFilename: <ipython-input-1-af507ad7b0d2>\n\nLine # Mem usage Increment Occurrences Line Contents\n=============================================================\n 12 135.8 MiB 135.8 MiB 1 @pandas_udf(\"long\")\n 13 def multiply_pandas_func(a: pd.Series, b: pd.Series) -> pd.Series:\n 14 143.5 MiB 7.8 MiB 1 return a * b\n\n\n============================================================\nProfile of UDF<id=2>\n============================================================\nFilename: <ipython-input-1-0e8713529486>\n\nLine # Mem usage Increment Occurrences Line Contents\n=============================================================\n 11 71.7 MiB 71.7 MiB 1 @arrow_udf(\"long\")\n 12 def multiply_arrow_func(a: pa.Array, b: pa.Array) -> pa.Array:\n 13 79.9 MiB 8.2 MiB 1 return pa.compute.multiply(a, b)ConclusioneDatabricks Runtime 18.0 introduce le UDF Arrow native, offrendo un'alternativa pi\u00f9 veloce e snella alle UDF Pandas per un'esecuzione performante delle UDF Python in PySpark. Operando direttamente sui dati Arrow ed eliminando l'overhead di conversione Pandas/Arrow, le UDF Arrow offrono un'esecuzione circa il 10% pi\u00f9 veloce, circa il 40% in meno di utilizzo di memoria e un migliore supporto per i tipi di dati complessi, il tutto con una sintassi a decoratore familiare e intuitiva.Pronto a esplorare di pi\u00f9? Prova oggi stesso le UDF Arrow native su Databricks come parte di Databricks Runtime 18.0. Per iniziare, sostituisci semplicemente le tue UDF Pandas esistenti con le UDF Arrow. Nella maggior parte dei casi, bastano poche righe di codice per sbloccare guadagni di prestazioni immediati. Consulta la documentazione delle UDF Arrow e la documentazione delle UDTF Arrow per il riferimento API completo e ulteriori esempi.(Questo post sul blog \u00e8 stato tradotto utilizzando strumenti basati sull'intelligenza artificiale) Post originale", "headline": "Introduzione alle Arrow UDFs in PySpark: Un sostituto pi\u00f9 veloce e leggero per le Pandas UDFs", "datePublished": " 05/20/2026T00:00:00-08:00", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_image_ImageObject", "url": "https://www.databricks.com/sites/default/files/2026-05/2026-02-blog-introducing-arrow-udfs-in-pyspark-inline-960x502.4.png"}], "dateModified": "05/22/2026T00:00:00-08:00", "mentions": [{"@type": "BreadcrumbList", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Tutti i blog", "item": "https://www.databricks.com/it/blog", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Ingegneria", "item": "https://www.databricks.com/it/blog/category/engineering", "position": 2}]}, {"name": "Python", "@id": "https://entity.schemaapp.com/DatabricksInc/CreativeWork_python_9ccca1ec3e5c4f56bb4441c4a6310e73b636329f3426d369a4b44d5ead8fdc42", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/05z1_", "http://www.wikidata.org/entity/Q28865", "https://en.wikipedia.org/wiki/Python_(programming_language)"]}, {"name": "Apache Spark", "@id": "https://entity.schemaapp.com/DatabricksInc/CreativeWork_apachespark_5c8dd45dbca4e0e307dfd60fd118717da7b8685e7a2e0f8c1c76c42d24a27cf6", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/0ndhxqz", "http://www.wikidata.org/entity/Q7573619", "https://en.wikipedia.org/wiki/Apache_Spark"]}, {"name": "Arrow missile", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_arrowmissile_f38b54b11bc27ae962ed682ee53261da2f19552ff691d2f5541efd50d3831f8c", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/01c5l9", "http://www.wikidata.org/entity/Q671836", "https://en.wikipedia.org/wiki/Arrow_(missile_family)"]}, {"name": "Arrow", "@id": "https://entity.schemaapp.com/DatabricksInc/CreativeWork_arrow_ccc427c842cf70c383ce0a6fa97e9fab844b3ee5c0bb90617c3378bef4e1a107", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/0l170__", "http://www.wikidata.org/entity/Q552314", "https://en.wikipedia.org/wiki/Arrow_(TV_series)"]}, {"name": "Possibile", "@id": "https://entity.schemaapp.com/DatabricksInc/Organization_possibile_881f8a0c0658172d96f88b9142153e240a5258daae0b3757931fb76d07d426f7", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/g/11btv7wkzt", "http://www.wikidata.org/entity/Q20970716", "https://en.wikipedia.org/wiki/Possible_(political_party)"]}, {"name": "note", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_note_e0501093bd0dec57cf6f5c18ed9cd63311096466daca23aa08414db86ff75d99", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/01wvdr", "http://www.wikidata.org/entity/Q12823770", "https://en.wikipedia.org/wiki/Note_(typography)"]}, {"name": "Indiana", "@id": "https://entity.schemaapp.com/DatabricksInc/Place_indiana_963f5ccc9a6be110cbe772a79e32de662d3b5e96e8248b734b42cfc65afdd0a3", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/03v1s", "http://www.wikidata.org/entity/Q1415", "https://en.wikipedia.org/wiki/Indiana"]}, {"name": "chief executive officer", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_chiefexecutiveofficer_3c3048b68a8a85ded766bf8e01673cd008ce2eb6f1766d180996104e8d39fe11", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/0dq_5", "http://www.wikidata.org/entity/Q484876", "https://en.wikipedia.org/wiki/Chief_executive_officer"]}], "articleSection": "Accedi\nOpen source", "author": [{"@type": "Person", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/it/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/it/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/it/blog/author/yicong-huang", "name": "Yicong Huang"}]}, {"@context": "http://schema.org", "@type": "Organization", "description": "The Databricks Platform is the world\u2019s first data intelligence platform powered by generative AI. Infuse AI into every facet of your business.", "name": "Databricks", "disambiguatingDescription": "Your data. Your AI. Your future. Own them all on the new data intelligence platform", "sameAs": ["kg:/m/0120wgnc", "https://www.wikidata.org/wiki/Q18350420", "https://en.wikipedia.org/wiki/Databricks", "https://twitter.com/databricks", "https://www.databricks.com/feed", "https://www.youtube.com/c/Databricks", "https://www.facebook.com/pages/Databricks/560203607379694", "https://www.linkedin.com/company/databricks", "https://www.glassdoor.com/Overview/Working-at-Databricks-EI_IE954734.11,21.htm"], "telephone": "+1-866-330-0121", "areaServed": "http://www.wikidata.org/entity/Q13780930", "legalName": "Databricks Inc.", "knowsLanguage": "en-US", "url": "https://www.databricks.com/", "additionalType": "https://www.wikidata.org/wiki/Q110029326", "logo": {"@type": "ImageObject", "width": "127", "height": "20", "url": "https://www.databricks.com/en-website-assets/static/8ed15a13c1511a75a4855999a2011c5c/f2f26/databricks-default.webp", "@id": "https://www.databricks.com/en-website-assets/static/8ed15a13c1511a75a4855999a2011c5c/f2f26/databricks-default.webp"}, "contactPoint": {"@type": "ContactPoint", "contactOption": "TollFree", "availableLanguage": "en-US", "contactType": "Contact", "telephone": "+1-866-330-0121", "name": "Databricks Contact", "@id": "https://www.databricks.com/#ContactPoint"}, "address": {"@type": "PostalAddress", "streetAddress": "160 Spear Street", "postalCode": "94105", "addressRegion": ["https://www.wikidata.org/wiki/Q99", "California"], "addressLocality": ["https://www.wikidata.org/wiki/Q62", "San Francisco"], "addressCountry": "http://www.wikidata.org/entity/Q30", "name": "Databricks Address", "@id": "https://www.databricks.com/#PostalAddress"}, "knowsAbout": [{"@type": "Thing", "sameAs": ["kg:/g/11khkg2rwf", "https://www.wikidata.org/wiki/Q117246174", "https://en.wikipedia.org/wiki/Generative_artificial_intelligence"], "name": "Generative AI", "@id": "https://www.databricks.com/#Thing"}, {"@type": "Thing", "sameAs": ["kg:/m/038_34", "https://en.wikipedia.org/wiki/Data_management", "https://www.wikidata.org/wiki/Q1149776"], "name": "data management", "@id": "https://www.databricks.com/#Thing1"}, {"@type": "Thing", "sameAs": ["kg:/m/0136zzks", "https://en.wikipedia.org/wiki/Data_lake", "https://www.wikidata.org/wiki/Q20707560"], "name": "data lake", "@id": "https://www.databricks.com/#Thing2"}, {"@type": "Thing", "sameAs": ["kg:/m/0h7m73m", "https://www.wikidata.org/wiki/Q2499178", "https://en.wikipedia.org/wiki/Cloud_database"], "name": "cloud database", "@id": "https://www.databricks.com/#Thing3"}, {"@type": "Thing", "sameAs": ["kg:/m/0mkz", "https://en.wikipedia.org/wiki/Artificial_intelligence", "https://www.wikidata.org/wiki/Q11660"], "name": "artificial intelligence", "@id": "https://www.databricks.com/#Thing4"}, {"@type": "Thing", "sameAs": ["kg:/m/01hyh", "https://www.wikidata.org/wiki/Q2539", "https://en.wikipedia.org/wiki/Machine_learning"], "name": "machine learning", "@id": "https://www.databricks.com/#Thing5"}], "image": {"@type": "ImageObject", "width": "1200", "height": "628", "url": "https://www.databricks.com/sites/default/files/2023-11/databricks-og-universal.png", "@id": "https://www.databricks.com/sites/default/files/2023-11/databricks-og-universal.png"}, "@id": "https://www.databricks.com/#Organization"}]
Passa al contenuto principale
Open source

Introduzione alle Arrow UDFs in PySpark: Un sostituto più veloce e leggero per le Pandas UDFs

Definisci UDF più performanti con facilità.

di Ruifeng Zheng e Yicong Huang

  • Introduciamo le UDF Arrow native, che operano direttamente sui dati Arrow, eliminando l'overhead di conversione Pandas/Arrow nelle UDF Pandas per un'esecuzione più rapida e un minore utilizzo di memoria.
  • Descriviamo anche i tipi di UDF Arrow per casi d'uso scalari e di aggregazione, e le UDTF Arrow per trasformazioni table-in, table-out, con esempi di codice sia in Python che in SQL.
  • I benchmark mostrano che le UDF Arrow sono circa il 10% più veloci e utilizzano circa il 40% in meno di memoria rispetto alle UDF Pandas, con un migliore supporto per i tipi di dati complessi.

Introduzione

Le funzioni definite dall'utente (UDF) Python sono un meccanismo di estensibilità essenziale, ma tradizionalmente hanno sofferto di un elevato overhead dovuto all'esecuzione basata su righe. In Apache Spark™, le UDF Pandas hanno affrontato parte di questo problema introducendo la serializzazione basata su Arrow e l'elaborazione batch, migliorando significativamente il throughput rispetto alle UDF Python scalari.

Tuttavia, le UDF Pandas presentano ancora limitazioni fondamentali:

  • La conversione dei dati Pandas/Arrow introduce copie di dati aggiuntive. Gli approcci a copia zero sono possibili solo in alcuni casi specifici. Ad esempio, le colonne con valori NULL attiveranno copie profonde.
  • I tipi di dati complessi non sono ben supportati. Ad esempio, le istanze di StructType nidificate non sono supportate per il tipo di output con casi d'uso di aggregazione.
Flussi di dati dell'esecuzione delle UDF Pandas in Apache Spark

Eliminando la conversione dei dati Pandas/Arrow, le UDF Arrow vengono eseguite più velocemente delle UDF Pandas, consumano meno memoria e offrono un migliore supporto per i tipi di dati.

UDF Arrow Native

Siamo entusiasti di presentare le UDF Arrow Native a partire da Databricks Runtime 18.0 (note di rilascio), un entusiasmante passo avanti per l'esecuzione performante delle UDF.

Le UDF Arrow Native operano direttamente sui dati Arrow senza convertire gli input in oggetti Pandas o NumPy. Ciò preserva il layout colonnare end-to-end, evita copie di dati non necessarie e consente alle UDF di utilizzare l'elaborazione vettorializzata sfruttando il modello di calcolo e memoria nativo di Arrow.

Per definire una UDF Arrow, gli utenti possono utilizzare un nuovo decoratore Python @arrow_udf, con tipo di ritorno specificato e tipo di valutazione opzionale. Ad esempio:

Gli utenti possono anche definirla con il decoratore esistente @udf con suggerimenti di tipo completi. Ad esempio:

Nota: la definizione della funzione dovrebbe includere suggerimenti di tipo per tutti gli argomenti e il valore di ritorno.
Questo design si allinea con le interfacce delle UDF Python scalari, fornendo un'esperienza coerente e intuitiva per gli utenti già familiari con le UDF Python scalari.

Quanto segue dimostra come utilizzare la UDF Arrow:

Utilizzo Python:

Utilizzo SQL:

Forniamo supporto per varianti delle interfacce UDF Arrow. Incluse funzioni scalari, funzioni di aggregazione e funzioni di tabella. Nell'API del dataframe forniamo anche mapInArrow e applyInArrow per utilizzare le UDF Arrow. Le introdurremo una per una.

Funzioni Scalari Arrow

Le funzioni scalari Arrow eseguono trasformazioni riga per riga. Sono l'equivalente Arrow delle UDF Pandas scalari e possono essere utilizzate ovunque sia prevista un'espressione di colonna, come df.select() o df.withColumn(). Sono supportate tre modalità di input: diretta, iteratore e iteratore di array multipli. Le varianti dell'iteratore sono utili quando la UDF richiede un'inizializzazione costosa una tantum (ad esempio, il caricamento di un modello o la compilazione di un pattern regex), poiché il costo di configurazione viene ammortizzato su tutti i batch. In tutti i casi, il numero di righe di output deve corrispondere al numero di righe di input.

  • Da Array ad Array: riceve uno o più pyarrow.Array e restituisce un pyarrow.Array. L'array di input e di output deve avere lo stesso numero di valori.
  • Da Iteratore di Array a Iteratore di Array: riceve un iteratore di pyarrow.Array e restituisce un iteratore di pyarrow.Array. Questo tipo è utile quando l'esecuzione della UDF richiede un'inizializzazione costosa.
  • Da Iteratore di Array Multipli a Iteratore di Array: riceve un iteratore di una tupla di più pyarrow.Array e restituisce un iteratore di pyarrow.Array.

Funzioni di Aggregazione Arrow

Le funzioni di aggregazione Arrow accettano uno o più input pyarrow.Array e restituiscono un valore scalare, riducendo un gruppo di righe in un singolo risultato. Sono l'equivalente Arrow delle UDF Pandas aggregate raggruppate e vengono utilizzate con groupBy().agg() o operazioni Window. Similmente alle funzioni scalari, le funzioni di aggregazione supportano anche tre modalità di input.

Da Array a Scalare: riceve pyarrow.Array e restituisce un valore scalare.

  • Da Iteratore di Array a Scalare: riceve un iteratore di pyarrow.Array e restituisce un valore scalare. Questo è utile per elaborare grandi volumi di dati in operazioni di aggregazione.

Da Iteratore di Array Multipli a Scalare: riceve un iteratore di una tupla di più pyarrow.Array e restituisce un valore scalare. È possibile definire aggregazioni più complesse.

Funzioni di Tabella Arrow

Le funzioni di tabella Arrow, note anche come UDTF Arrow (User-Defined Table Functions), accettano un pyarrow.RecordBatch o più pa.Array come input e producono un pyarrow.Table come output. Questo rappresenta il modello predominante per le trasformazioni table-in, table-out implementate in Python utilizzando l'esecuzione colonnare. Le UDTF Arrow possiedono la capacità di:

  • Restituire più colonne
  • Produrre zero, una o più righe
  • Eseguire trasformazioni di tabella vettorializzate impiegando kernel di calcolo Arrow

Di conseguenza, sono ottimamente adatte per operazioni come il filtraggio, l'espansione delle righe, la ristrutturazione dei dati e la generazione di colonne derivate.

L'interfaccia arrow_udtf è progettata per la semplicità, impiegando una sintassi a decoratore in cui si definisce il tipo di ritorno usando una stringa formattata DDL. In questa configurazione, il metodo eval accetta oggetti PyArrow come input e si prevede che produca tabelle PyArrow o RecordBatches. L'interfaccia supporta due modalità di input. Quando si elaborano argomenti di tabella, al metodo eval viene fornito un oggetto pa.RecordBatch che incapsula tutte le colonne dalla tabella di input:

Per gli argomenti scalari, il metodo riceve oggetti pa.Array, uno per ogni input scalare:

Ecco un altro esempio:

Questa UDTF può funzionare in due modi distinti:

Utilizzo Python:

Utilizzo SQL:

Supporto per DataFrame mapInArrow e applyInArrow

Oltre alle User-Defined Functions (UDF) e alle User-Defined Table Functions (UDTF), PySpark fornisce API di funzione Arrow che facilitano l'applicazione diretta di funzioni native Python ai dati Arrow a livello di DataFrame. Queste API operano in modo analogo alle loro controparti Pandas (mapInPandas, applyInPandas) ma utilizzano pyarrow.RecordBatch e pyarrow.Table invece dei DataFrame Pandas, aggirando così l'overhead di conversione tra i formati Pandas e Arrow.

  • Mappa. DataFrame.mapInArrow trasforma un iteratore di pyarrow.RecordBatch in un altro iteratore di pyarrow.RecordBatch, consentendo operazioni a livello di riga come il filtraggio, la trasformazione o l'espansione.
  • Mappa raggruppata. groupBy().applyInArrow() applica una funzione specificata a ogni gruppo, accettando e restituendo un pyarrow.Table. Questa funzionalità si rivela utile per le trasformazioni per gruppo, come la normalizzazione dei dati.
  • Mappa co-raggruppata. cogroup().applyInArrow() consente il co-raggruppamento di due DataFrame basati su una chiave condivisa, applicando successivamente una funzione a ogni co-gruppo. La funzione riceve due pyarrow.Table input e si prevede che restituisca un singolo pyarrow.Table.

Prestazioni

Eliminando la costosa conversione dei dati Pandas/Arrow, le UDF Arrow generalmente vengono eseguite più velocemente delle UDF Pandas, con un minore utilizzo di memoria. Confrontiamo le due semplici UDF:

L'UDF Arrow è circa il 10% più veloce dell'UDF Pandas, e il profiler di memoria mostra che circa il 40% di memoria viene risparmiato nell'esecuzione.

Conclusione

Databricks Runtime 18.0 introduce le UDF Arrow native, offrendo un'alternativa più veloce e snella alle UDF Pandas per un'esecuzione performante delle UDF Python in PySpark. Operando direttamente sui dati Arrow ed eliminando l'overhead di conversione Pandas/Arrow, le UDF Arrow offrono un'esecuzione circa il 10% più veloce, circa il 40% in meno di utilizzo di memoria e un migliore supporto per i tipi di dati complessi, il tutto con una sintassi a decoratore familiare e intuitiva.

Pronto a esplorare di più? Prova oggi stesso le UDF Arrow native su Databricks come parte di Databricks Runtime 18.0. Per iniziare, sostituisci semplicemente le tue UDF Pandas esistenti con le UDF Arrow. Nella maggior parte dei casi, bastano poche righe di codice per sbloccare guadagni di prestazioni immediati. Consulta la documentazione delle UDF Arrow e la documentazione delle UDTF Arrow per il riferimento API completo e ulteriori esempi.

(Questo post sul blog è stato tradotto utilizzando strumenti basati sull'intelligenza artificiale) Post originale

Ricevi gli ultimi articoli nella tua casella di posta

Iscriviti al nostro blog e ricevi gli ultimi articoli direttamente nella tua casella di posta.