Direkt zum Hauptinhalt
Open-Source

Einführung von Arrow UDFs in PySpark: Ein schnellerer, schlankerer Ersatz für Pandas UDFs

Definieren Sie performantere UDFs mit Leichtigkeit.

von Ruifeng Zheng und Yicong Huang

  • Wir stellen native Arrow UDFs vor, die direkt mit Arrow-Daten arbeiten und den Pandas/Arrow-Konvertierungs-Overhead in Pandas UDFs eliminieren, was zu schnellerer Ausführung und geringerem Speicherverbrauch führt.
  • Wir beschreiben auch Arrow UDF-Typen für Skalar- und Aggregations-Anwendungsfälle sowie Arrow UDTFs für Table-in-, Table-out-Transformationen, mit Codebeispielen in Python und SQL.
  • Benchmarks zeigen, dass Arrow UDFs ~10% schneller sind und ~40% weniger Speicher als Pandas UDFs verbrauchen, mit besserer Unterstützung für komplexe Datentypen.

Einführung

Python benutzerdefinierte Funktionen (UDFs) sind ein wesentlicher Erweiterungsmechanismus, litten jedoch traditionell unter hohem Overhead aufgrund der zeilenbasierten Ausführung. In Apache Spark™ lösten Pandas UDFs einen Teil dieses Problems durch die Einführung von Arrow-basierter Serialisierung und Stapelverarbeitung, die den Durchsatz im Vergleich zu skalaren Python UDFs erheblich verbesserten.

Pandas UDFs weisen jedoch immer noch grundlegende Einschränkungen auf:

  • Die Pandas/Arrow-Datenkonvertierung führt zu zusätzlichen Datenkopien. Zero-Copy-Ansätze sind nur in bestimmten engen Fällen möglich. Zum Beispiel lösen Spalten mit NULL-Werten tiefe Kopien aus.
  • Komplexe Datentypen werden nicht gut unterstützt. Zum Beispiel werden verschachtelte StructType-Instanzen für den Ausgabetyp bei Aggregationsanwendungsfällen nicht unterstützt.
Datenflüsse der Pandas UDF-Ausführung in Apache Spark

Durch das Weglassen der Pandas/Arrow-Datenkonvertierung führen die Arrow UDFs schneller aus als Pandas UDFs, verbrauchen weniger Speicher und bieten eine bessere Datentypunterstützung.

Native Arrow UDFs

Wir freuen uns, die Einführung von Native Arrow UDFs ab Databricks Runtime 18.0 (Versionshinweise) bekannt zu geben, ein aufregender Fortschritt für die performante UDF-Ausführung.

Native Arrow UDFs arbeiten direkt mit Arrow-Daten, ohne Eingaben in Pandas- oder NumPy-Objekte zu konvertieren. Dies bewahrt das spaltenbasierte Layout durchgängig, vermeidet unnötige Datenkopien und ermöglicht UDFs die Nutzung von vektorisierter Verarbeitung durch die Ausnutzung von Arrows nativem Berechnungs- und Speichermodell.

Um eine Arrow UDF zu definieren, können Benutzer einen neuen Python-Decorator @arrow_udf verwenden, mit einem angegebenen Rückgabetyp und einem optionalen Auswertungstyp. Zum Beispiel:

Benutzer können sie auch mit dem bestehenden Decorator @udf und vollständigen Typ-Hinweisen definieren. Zum Beispiel:

Hinweis: Die Funktionsdefinition sollte Typ-Hinweise für alle Argumente und den Rückgabewert enthalten.
Dieses Design stimmt mit den Schnittstellen skalarer Python UDFs überein und bietet eine konsistente und intuitive Erfahrung für Benutzer, die bereits mit skalaren Python UDFs vertraut sind.

Das Folgende zeigt, wie man die Arrow UDF verwendet:

Python-Nutzung:

SQL-Nutzung:

Wir bieten Unterstützung für Varianten von Arrow UDF-Schnittstellen. Dazu gehören Skalarfunktionen, Aggregatfunktionen und Tabellenfunktionen. In der DataFrame API stellen wir auch mapInArrow und applyInArrow zur Verfügung, um Arrow UDFs zu verwenden. Wir werden sie im Folgenden einzeln vorstellen.

Arrow Skalarfunktionen

Arrow Skalarfunktionen führen zeilenweise Transformationen durch. Sie sind das Arrow-Äquivalent von skalaren Pandas UDFs und können überall dort verwendet werden, wo ein Spaltenausdruck erwartet wird, wie zum Beispiel df.select() oder df.withColumn(). Es werden drei Eingabemodi unterstützt: direkt, Iterator und Iterator mehrerer Arrays. Die Iterator-Varianten sind nützlich, wenn die UDF eine aufwendige einmalige Initialisierung erfordert (z. B. das Laden eines Modells oder das Kompilieren eines Regex-Musters), da die Einrichtungskosten über alle Batches amortisiert werden. In allen Fällen muss die Anzahl der Ausgabezellen mit der Anzahl der Eingabezellen übereinstimmen.

  • Arrays zu Array: Empfängt ein oder mehrere pyarrow.Array und gibt ein pyarrow.Array zurück. Das Eingabe- und Ausgabe-Array muss die gleiche Anzahl von Werten haben.
  • Iterator von Arrays zu Iterator von Arrays: Empfängt einen Iterator von pyarrow.Array und gibt einen Iterator von pyarrow.Array zurück. Dieser Typ ist nützlich, wenn die UDF-Ausführung eine aufwendige Initialisierung erfordert.
  • Iterator mehrerer Arrays zu Iterator von Arrays: Empfängt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Iterator von pyarrow.Array zurück.

Arrow Aggregatfunktionen

Arrow Aggregatfunktionen nehmen eine oder mehrere pyarrow.Array-Eingaben entgegen und geben einen Skalarwert zurück, wobei eine Gruppe von Zeilen zu einem einzigen Ergebnis reduziert wird. Sie sind das Arrow-Äquivalent von gruppierten Aggregat-Pandas UDFs und werden mit groupBy().agg() oder Window-Operationen verwendet. Ähnlich wie Skalarfunktionen unterstützen Aggregatfunktionen ebenfalls drei Eingabemodi.

Arrays zu Skalar: Empfängt pyarrow.Array und gibt einen Skalarwert zurück.

  • Iterator von Arrays zu Skalar: Empfängt einen Iterator von pyarrow.Array und gibt einen Skalarwert zurück. Dies ist nützlich für die Verarbeitung großer Datenmengen in aggregationsartigen Operationen.

Iterator mehrerer Arrays zu Skalar: Empfängt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Skalarwert zurück. Komplexere Aggregationen können definiert werden.

Arrow Tabellenfunktionen

Arrow Tabellenfunktionen, auch bekannt als Arrow UDTFs (benutzerdefinierte Tabellenfunktionen), akzeptieren ein pyarrow.RecordBatch oder mehrere pa.Array als Eingabe und erzeugen ein pyarrow.Table als Ausgabe. Dies stellt das vorherrschende Muster für Tabellen-in, Tabellen-out-Transformationen dar, die in Python unter Verwendung von spaltenbasierter Ausführung implementiert werden. Arrow UDTFs bieten die Möglichkeit,:

  • Mehrere Spalten zurückgeben
  • Null, eine oder mehrere Zeilen erzeugen
  • Vektorisierte Tabellentransformationen unter Verwendung von Arrow Compute-Kernels ausführen

Folglich sind sie optimal für Operationen wie Filtern, Zeilenerweiterung, Datenrestrukturierung und die Generierung abgeleiteter Spalten geeignet.

Die arrow_udtf Schnittstelle ist auf Einfachheit ausgelegt und verwendet eine Decorator-Syntax, bei der Sie den Rückgabetyp mithilfe eines DDL-formatierten Strings definieren. In diesem Setup akzeptiert die eval Methode PyArrow-Objekte als Eingabe und soll PyArrow Tables oder RecordBatches liefern. Die Schnittstelle unterstützt zwei Eingabemodi. Bei der Verarbeitung von Tabellenargumenten erhält die eval Methode ein pa.RecordBatch Objekt, das alle Spalten der Eingabetabelle kapselt:

Für skalare Argumente empfängt die Methode pa.Array-Objekte, eines für jede skalare Eingabe:

Hier ist ein weiteres Beispiel:

Diese UDTF kann auf zwei verschiedene Arten funktionieren:

Python-Nutzung:

SQL-Nutzung:

Unterstützung für DataFrame mapInArrow und applyInArrow

Zusätzlich zu benutzerdefinierten Funktionen (UDFs) und benutzerdefinierten Tabellenfunktionen (UDTFs) bietet PySpark Arrow Function APIs, die die direkte Anwendung nativer Python-Funktionen auf Arrow-Daten auf DataFrame-Ebene ermöglichen. Diese APIs funktionieren analog zu ihren Pandas-Pendants (mapInPandas, applyInPandas) verwenden jedoch pyarrow.RecordBatch und pyarrow.Table anstelle von Pandas DataFrames, wodurch der Konvertierungsaufwand zwischen Pandas- und Arrow-Formaten umgangen wird.

  • Map. DataFrame.mapInArrow transformiert einen Iterator von pyarrow.RecordBatch in einen anderen Iterator von pyarrow.RecordBatch, wodurch zeilenbasierte Operationen wie Filtern, Transformation oder Erweiterung ermöglicht werden.
  • Grouped Map. groupBy().applyInArrow() wendet eine angegebene Funktion auf jede Gruppe an, die ein pyarrow.Table akzeptiert und zurückgibt. Diese Funktionalität erweist sich als vorteilhaft für gruppenweise Transformationen, wie z.B. Daten-Normalisierung.
  • Co-grouped Map. cogroup().applyInArrow() ermöglicht das Cogrouping von zwei DataFrames basierend auf einem gemeinsamen Schlüssel und wendet anschließend eine Funktion auf jede Cogroup an. Die Funktion empfängt zwei pyarrow.Table Eingaben und soll eine einzelne pyarrow.Table zurückgeben.

Leistung

Durch das Entfernen der aufwendigen Pandas/Arrow-Datenkonvertierung werden Arrow UDFs im Allgemeinen schneller ausgeführt als Pandas UDFs, mit geringerem Speicherverbrauch. Vergleichen wir die beiden einfachen UDFs:

Die Arrow UDF ist ~10% schneller als die Pandas UDF, und der Speicherprofiler zeigt, dass ~40% Speicher bei der Ausführung eingespart werden.

Fazit

Databricks Runtime 18.0 führt native Arrow UDFs ein, die eine schnellere, schlankere Alternative zu Pandas UDFs für die performante Ausführung von Python UDFs in PySpark bieten. Durch die direkte Verarbeitung von Arrow-Daten und die Eliminierung des Pandas/Arrow-Konvertierungsaufwands ermöglichen Arrow UDFs eine ~10% schnellere Ausführung, ~40% weniger Speicherverbrauch und eine bessere Unterstützung für komplexe Datentypen – alles mit einer vertrauten, intuitiven Decorator-Syntax.

Bereit, mehr zu entdecken? Probieren Sie Native Arrow UDFs noch heute auf Databricks als Teil von Databricks Runtime 18.0 aus. Um zu beginnen, ersetzen Sie einfach Ihre bestehenden Pandas UDFs durch Arrow UDFs. In den meisten Fällen sind nur wenige Zeilen Codeänderung erforderlich, um sofortige Leistungssteigerungen zu erzielen. Weitere Informationen finden Sie in der Arrow UDF-Dokumentation und der Arrow UDTF-Dokumentation für die vollständige API-Referenz und zusätzliche Beispiele.

(Dieser Blogbeitrag wurde mit KI-gestützten Tools übersetzt.) Originalbeitrag

Erhalten Sie die neuesten Beiträge in Ihrem Posteingang

Abonnieren Sie unseren Blog und erhalten Sie die neuesten Beiträge direkt in Ihren Posteingang.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])Iterator mehrerer Arrays zu Iterator von Arrays: Empf\u00e4ngt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Iterator von pyarrow.Array zur\u00fcck.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)Arrow AggregatfunktionenArrow Aggregatfunktionen nehmen eine oder mehrere pyarrow.Array-Eingaben entgegen und geben einen Skalarwert zur\u00fcck, wobei eine Gruppe von Zeilen zu einem einzigen Ergebnis reduziert wird. Sie sind das Arrow-\u00c4quivalent von gruppierten Aggregat-Pandas UDFs und werden mit groupBy().agg() oder Window-Operationen verwendet. \u00c4hnlich wie Skalarfunktionen unterst\u00fctzen Aggregatfunktionen ebenfalls drei Eingabemodi. Arrays zu Skalar: Empf\u00e4ngt pyarrow.Array und gibt einen Skalarwert zur\u00fcck. 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)Iterator von Arrays zu Skalar: Empf\u00e4ngt einen Iterator von pyarrow.Array und gibt einen Skalarwert zur\u00fcck. Dies ist n\u00fctzlich f\u00fcr die Verarbeitung gro\u00dfer Datenmengen in aggregationsartigen Operationen.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.0Iterator mehrerer Arrays zu Skalar: Empf\u00e4ngt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Skalarwert zur\u00fcck. Komplexere Aggregationen k\u00f6nnen definiert werden.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.0Arrow TabellenfunktionenArrow Tabellenfunktionen, auch bekannt als Arrow UDTFs (benutzerdefinierte Tabellenfunktionen), akzeptieren ein pyarrow.RecordBatch oder mehrere pa.Array als Eingabe und erzeugen ein pyarrow.Table als Ausgabe. Dies stellt das vorherrschende Muster f\u00fcr Tabellen-in, Tabellen-out-Transformationen dar, die in Python unter Verwendung von spaltenbasierter Ausf\u00fchrung implementiert werden. Arrow UDTFs bieten die M\u00f6glichkeit,:Mehrere Spalten zur\u00fcckgebenNull, eine oder mehrere Zeilen erzeugenVektorisierte Tabellentransformationen unter Verwendung von Arrow Compute-Kernels ausf\u00fchrenFolglich sind sie optimal f\u00fcr Operationen wie Filtern, Zeilenerweiterung, Datenrestrukturierung und die Generierung abgeleiteter Spalten geeignet.Die arrow_udtf Schnittstelle ist auf Einfachheit ausgelegt und verwendet eine Decorator-Syntax, bei der Sie den R\u00fcckgabetyp mithilfe eines DDL-formatierten Strings definieren. In diesem Setup akzeptiert die eval Methode PyArrow-Objekte als Eingabe und soll PyArrow Tables oder RecordBatches liefern. Die Schnittstelle unterst\u00fctzt zwei Eingabemodi. Bei der Verarbeitung von Tabellenargumenten erh\u00e4lt die eval Methode ein pa.RecordBatch Objekt, das alle Spalten der Eingabetabelle kapselt: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 resultF\u00fcr skalare Argumente empf\u00e4ngt die Methode pa.Array-Objekte, eines f\u00fcr jede skalare Eingabe: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 resultHier ist ein weiteres Beispiel: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 })Diese UDTF kann auf zwei verschiedene Arten funktionieren:Python-Nutzung: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# +----------+-------+--------+SQL-Nutzung: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# +----------+-------+--------+Unterst\u00fctzung f\u00fcr DataFrame mapInArrow und applyInArrowZus\u00e4tzlich zu benutzerdefinierten Funktionen (UDFs) und benutzerdefinierten Tabellenfunktionen (UDTFs) bietet PySpark Arrow Function APIs, die die direkte Anwendung nativer Python-Funktionen auf Arrow-Daten auf DataFrame-Ebene erm\u00f6glichen. Diese APIs funktionieren analog zu ihren Pandas-Pendants (mapInPandas, applyInPandas) verwenden jedoch pyarrow.RecordBatch und pyarrow.Table anstelle von Pandas DataFrames, wodurch der Konvertierungsaufwand zwischen Pandas- und Arrow-Formaten umgangen wird.Map. DataFrame.mapInArrow transformiert einen Iterator von pyarrow.RecordBatch in einen anderen Iterator von pyarrow.RecordBatch, wodurch zeilenbasierte Operationen wie Filtern, Transformation oder Erweiterung erm\u00f6glicht werden.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# +---+---+Grouped Map. groupBy().applyInArrow() wendet eine angegebene Funktion auf jede Gruppe an, die ein pyarrow.Table akzeptiert und zur\u00fcckgibt. Diese Funktionalit\u00e4t erweist sich als vorteilhaft f\u00fcr gruppenweise Transformationen, wie z.B. Daten-Normalisierung.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# +---+-------------------+Co-grouped Map. cogroup().applyInArrow() erm\u00f6glicht das Cogrouping von zwei DataFrames basierend auf einem gemeinsamen Schl\u00fcssel und wendet anschlie\u00dfend eine Funktion auf jede Cogroup an. Die Funktion empf\u00e4ngt zwei pyarrow.Table Eingaben und soll eine einzelne pyarrow.Table zur\u00fcckgeben.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# +----+-----+LeistungDurch das Entfernen der aufwendigen Pandas/Arrow-Datenkonvertierung werden Arrow UDFs im Allgemeinen schneller ausgef\u00fchrt als Pandas UDFs, mit geringerem Speicherverbrauch. Vergleichen wir die beiden einfachen UDFs: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)Die Arrow UDF ist ~10% schneller als die Pandas UDF, und der Speicherprofiler zeigt, dass ~40% Speicher bei der Ausf\u00fchrung eingespart werden.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)FazitDatabricks Runtime 18.0 f\u00fchrt native Arrow UDFs ein, die eine schnellere, schlankere Alternative zu Pandas UDFs f\u00fcr die performante Ausf\u00fchrung von Python UDFs in PySpark bieten. Durch die direkte Verarbeitung von Arrow-Daten und die Eliminierung des Pandas/Arrow-Konvertierungsaufwands erm\u00f6glichen Arrow UDFs eine ~10% schnellere Ausf\u00fchrung, ~40% weniger Speicherverbrauch und eine bessere Unterst\u00fctzung f\u00fcr komplexe Datentypen \u2013 alles mit einer vertrauten, intuitiven Decorator-Syntax.Bereit, mehr zu entdecken? Probieren Sie Native Arrow UDFs noch heute auf Databricks als Teil von Databricks Runtime 18.0 aus. Um zu beginnen, ersetzen Sie einfach Ihre bestehenden Pandas UDFs durch Arrow UDFs. In den meisten F\u00e4llen sind nur wenige Zeilen Code\u00e4nderung erforderlich, um sofortige Leistungssteigerungen zu erzielen. Weitere Informationen finden Sie in der Arrow UDF-Dokumentation und der Arrow UDTF-Dokumentation f\u00fcr die vollst\u00e4ndige API-Referenz und zus\u00e4tzliche Beispiele.(Dieser Blogbeitrag wurde mit KI-gest\u00fctzten Tools \u00fcbersetzt.) Originalbeitrag", "headline": "Einf\u00fchrung von Arrow UDFs in PySpark: Ein schnellerer, schlankerer Ersatz f\u00fcr Pandas UDFs", "mainEntityOfPage": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/de/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": "22.05.2026T00:00:00-08:00", "datePublished": " 20.05.2026T00:00:00-08:00", "articleSection": "Anmelden\nOpen-Source", "mentions": [{"@type": "BreadcrumbList", "@id": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Alle Blogs", "item": "https://www.databricks.com/de/blog", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Engineering", "item": "https://www.databricks.com/de/blog/category/engineering", "position": 2}]}, {"name": "serialization", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_serialization_3d635538277e0801ee547bad451b3df91592806ce9ecf24b3f2461c92feb324d", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q1127410"}, {"name": "batch processing", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_batchprocessing_e9bde7411996b86bf8a0dde825ba179d7e939b58880a3cbd55819bde61c89f1d", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q661613"}, {"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": "throughput", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_throughput_692f3814e812c5a5c30220df3c48aba30a1632e9480456525e8ec53e1f0eb714", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q1383412"}, {"name": "overhead", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_overhead_a93dc864bece221f04fa6ff23bb11d37d2123fe7c54af0ee3705b03e0e99459d", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q2006368"}, {"name": "data type", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_datatype_d259ed375d757ffa033beda2d4482ab4ee3bae3a820bb08bc9bb9302f5092757", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q190087"}, {"name": "value", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_value_7681edb4f4a2d9b519ad53b393b1304d14506b691302092f1c7e4f06ac3de976", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q194112"}, {"name": "function", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_function_15bf4498f958ebce51f3ee47cb5b4b24a3efda18e4dd44e45cfec18e0ba044e3", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": "http://www.wikidata.org/entity/Q11348"}, {"name": "die", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_die_f9eceedb2f860abfb5e8790fac57a8178ad90bb0c4ffc1d5d7d0150218e6beed", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/026bcl6", "http://www.wikidata.org/entity/Q1072430", "https://en.wikipedia.org/wiki/Die_(integrated_circuit)"]}, {"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"]}], "author": [{"@type": "Person", "@id": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/de/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/de/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/de/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"}]
Direkt zum Hauptinhalt
Open-Source

Einführung von Arrow UDFs in PySpark: Ein schnellerer, schlankerer Ersatz für Pandas UDFs

Definieren Sie performantere UDFs mit Leichtigkeit.

von Ruifeng Zheng und Yicong Huang

  • Wir stellen native Arrow UDFs vor, die direkt mit Arrow-Daten arbeiten und den Pandas/Arrow-Konvertierungs-Overhead in Pandas UDFs eliminieren, was zu schnellerer Ausführung und geringerem Speicherverbrauch führt.
  • Wir beschreiben auch Arrow UDF-Typen für Skalar- und Aggregations-Anwendungsfälle sowie Arrow UDTFs für Table-in-, Table-out-Transformationen, mit Codebeispielen in Python und SQL.
  • Benchmarks zeigen, dass Arrow UDFs ~10% schneller sind und ~40% weniger Speicher als Pandas UDFs verbrauchen, mit besserer Unterstützung für komplexe Datentypen.

Einführung

Python benutzerdefinierte Funktionen (UDFs) sind ein wesentlicher Erweiterungsmechanismus, litten jedoch traditionell unter hohem Overhead aufgrund der zeilenbasierten Ausführung. In Apache Spark™ lösten Pandas UDFs einen Teil dieses Problems durch die Einführung von Arrow-basierter Serialisierung und Stapelverarbeitung, die den Durchsatz im Vergleich zu skalaren Python UDFs erheblich verbesserten.

Pandas UDFs weisen jedoch immer noch grundlegende Einschränkungen auf:

  • Die Pandas/Arrow-Datenkonvertierung führt zu zusätzlichen Datenkopien. Zero-Copy-Ansätze sind nur in bestimmten engen Fällen möglich. Zum Beispiel lösen Spalten mit NULL-Werten tiefe Kopien aus.
  • Komplexe Datentypen werden nicht gut unterstützt. Zum Beispiel werden verschachtelte StructType-Instanzen für den Ausgabetyp bei Aggregationsanwendungsfällen nicht unterstützt.
Datenflüsse der Pandas UDF-Ausführung in Apache Spark

Durch das Weglassen der Pandas/Arrow-Datenkonvertierung führen die Arrow UDFs schneller aus als Pandas UDFs, verbrauchen weniger Speicher und bieten eine bessere Datentypunterstützung.

Native Arrow UDFs

Wir freuen uns, die Einführung von Native Arrow UDFs ab Databricks Runtime 18.0 (Versionshinweise) bekannt zu geben, ein aufregender Fortschritt für die performante UDF-Ausführung.

Native Arrow UDFs arbeiten direkt mit Arrow-Daten, ohne Eingaben in Pandas- oder NumPy-Objekte zu konvertieren. Dies bewahrt das spaltenbasierte Layout durchgängig, vermeidet unnötige Datenkopien und ermöglicht UDFs die Nutzung von vektorisierter Verarbeitung durch die Ausnutzung von Arrows nativem Berechnungs- und Speichermodell.

Um eine Arrow UDF zu definieren, können Benutzer einen neuen Python-Decorator @arrow_udf verwenden, mit einem angegebenen Rückgabetyp und einem optionalen Auswertungstyp. Zum Beispiel:

Benutzer können sie auch mit dem bestehenden Decorator @udf und vollständigen Typ-Hinweisen definieren. Zum Beispiel:

Hinweis: Die Funktionsdefinition sollte Typ-Hinweise für alle Argumente und den Rückgabewert enthalten.
Dieses Design stimmt mit den Schnittstellen skalarer Python UDFs überein und bietet eine konsistente und intuitive Erfahrung für Benutzer, die bereits mit skalaren Python UDFs vertraut sind.

Das Folgende zeigt, wie man die Arrow UDF verwendet:

Python-Nutzung:

SQL-Nutzung:

Wir bieten Unterstützung für Varianten von Arrow UDF-Schnittstellen. Dazu gehören Skalarfunktionen, Aggregatfunktionen und Tabellenfunktionen. In der DataFrame API stellen wir auch mapInArrow und applyInArrow zur Verfügung, um Arrow UDFs zu verwenden. Wir werden sie im Folgenden einzeln vorstellen.

Arrow Skalarfunktionen

Arrow Skalarfunktionen führen zeilenweise Transformationen durch. Sie sind das Arrow-Äquivalent von skalaren Pandas UDFs und können überall dort verwendet werden, wo ein Spaltenausdruck erwartet wird, wie zum Beispiel df.select() oder df.withColumn(). Es werden drei Eingabemodi unterstützt: direkt, Iterator und Iterator mehrerer Arrays. Die Iterator-Varianten sind nützlich, wenn die UDF eine aufwendige einmalige Initialisierung erfordert (z. B. das Laden eines Modells oder das Kompilieren eines Regex-Musters), da die Einrichtungskosten über alle Batches amortisiert werden. In allen Fällen muss die Anzahl der Ausgabezellen mit der Anzahl der Eingabezellen übereinstimmen.

  • Arrays zu Array: Empfängt ein oder mehrere pyarrow.Array und gibt ein pyarrow.Array zurück. Das Eingabe- und Ausgabe-Array muss die gleiche Anzahl von Werten haben.
  • Iterator von Arrays zu Iterator von Arrays: Empfängt einen Iterator von pyarrow.Array und gibt einen Iterator von pyarrow.Array zurück. Dieser Typ ist nützlich, wenn die UDF-Ausführung eine aufwendige Initialisierung erfordert.
  • Iterator mehrerer Arrays zu Iterator von Arrays: Empfängt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Iterator von pyarrow.Array zurück.

Arrow Aggregatfunktionen

Arrow Aggregatfunktionen nehmen eine oder mehrere pyarrow.Array-Eingaben entgegen und geben einen Skalarwert zurück, wobei eine Gruppe von Zeilen zu einem einzigen Ergebnis reduziert wird. Sie sind das Arrow-Äquivalent von gruppierten Aggregat-Pandas UDFs und werden mit groupBy().agg() oder Window-Operationen verwendet. Ähnlich wie Skalarfunktionen unterstützen Aggregatfunktionen ebenfalls drei Eingabemodi.

Arrays zu Skalar: Empfängt pyarrow.Array und gibt einen Skalarwert zurück.

  • Iterator von Arrays zu Skalar: Empfängt einen Iterator von pyarrow.Array und gibt einen Skalarwert zurück. Dies ist nützlich für die Verarbeitung großer Datenmengen in aggregationsartigen Operationen.

Iterator mehrerer Arrays zu Skalar: Empfängt einen Iterator eines Tupels mehrerer pyarrow.Array und gibt einen Skalarwert zurück. Komplexere Aggregationen können definiert werden.

Arrow Tabellenfunktionen

Arrow Tabellenfunktionen, auch bekannt als Arrow UDTFs (benutzerdefinierte Tabellenfunktionen), akzeptieren ein pyarrow.RecordBatch oder mehrere pa.Array als Eingabe und erzeugen ein pyarrow.Table als Ausgabe. Dies stellt das vorherrschende Muster für Tabellen-in, Tabellen-out-Transformationen dar, die in Python unter Verwendung von spaltenbasierter Ausführung implementiert werden. Arrow UDTFs bieten die Möglichkeit,:

  • Mehrere Spalten zurückgeben
  • Null, eine oder mehrere Zeilen erzeugen
  • Vektorisierte Tabellentransformationen unter Verwendung von Arrow Compute-Kernels ausführen

Folglich sind sie optimal für Operationen wie Filtern, Zeilenerweiterung, Datenrestrukturierung und die Generierung abgeleiteter Spalten geeignet.

Die arrow_udtf Schnittstelle ist auf Einfachheit ausgelegt und verwendet eine Decorator-Syntax, bei der Sie den Rückgabetyp mithilfe eines DDL-formatierten Strings definieren. In diesem Setup akzeptiert die eval Methode PyArrow-Objekte als Eingabe und soll PyArrow Tables oder RecordBatches liefern. Die Schnittstelle unterstützt zwei Eingabemodi. Bei der Verarbeitung von Tabellenargumenten erhält die eval Methode ein pa.RecordBatch Objekt, das alle Spalten der Eingabetabelle kapselt:

Für skalare Argumente empfängt die Methode pa.Array-Objekte, eines für jede skalare Eingabe:

Hier ist ein weiteres Beispiel:

Diese UDTF kann auf zwei verschiedene Arten funktionieren:

Python-Nutzung:

SQL-Nutzung:

Unterstützung für DataFrame mapInArrow und applyInArrow

Zusätzlich zu benutzerdefinierten Funktionen (UDFs) und benutzerdefinierten Tabellenfunktionen (UDTFs) bietet PySpark Arrow Function APIs, die die direkte Anwendung nativer Python-Funktionen auf Arrow-Daten auf DataFrame-Ebene ermöglichen. Diese APIs funktionieren analog zu ihren Pandas-Pendants (mapInPandas, applyInPandas) verwenden jedoch pyarrow.RecordBatch und pyarrow.Table anstelle von Pandas DataFrames, wodurch der Konvertierungsaufwand zwischen Pandas- und Arrow-Formaten umgangen wird.

  • Map. DataFrame.mapInArrow transformiert einen Iterator von pyarrow.RecordBatch in einen anderen Iterator von pyarrow.RecordBatch, wodurch zeilenbasierte Operationen wie Filtern, Transformation oder Erweiterung ermöglicht werden.
  • Grouped Map. groupBy().applyInArrow() wendet eine angegebene Funktion auf jede Gruppe an, die ein pyarrow.Table akzeptiert und zurückgibt. Diese Funktionalität erweist sich als vorteilhaft für gruppenweise Transformationen, wie z.B. Daten-Normalisierung.
  • Co-grouped Map. cogroup().applyInArrow() ermöglicht das Cogrouping von zwei DataFrames basierend auf einem gemeinsamen Schlüssel und wendet anschließend eine Funktion auf jede Cogroup an. Die Funktion empfängt zwei pyarrow.Table Eingaben und soll eine einzelne pyarrow.Table zurückgeben.

Leistung

Durch das Entfernen der aufwendigen Pandas/Arrow-Datenkonvertierung werden Arrow UDFs im Allgemeinen schneller ausgeführt als Pandas UDFs, mit geringerem Speicherverbrauch. Vergleichen wir die beiden einfachen UDFs:

Die Arrow UDF ist ~10% schneller als die Pandas UDF, und der Speicherprofiler zeigt, dass ~40% Speicher bei der Ausführung eingespart werden.

Fazit

Databricks Runtime 18.0 führt native Arrow UDFs ein, die eine schnellere, schlankere Alternative zu Pandas UDFs für die performante Ausführung von Python UDFs in PySpark bieten. Durch die direkte Verarbeitung von Arrow-Daten und die Eliminierung des Pandas/Arrow-Konvertierungsaufwands ermöglichen Arrow UDFs eine ~10% schnellere Ausführung, ~40% weniger Speicherverbrauch und eine bessere Unterstützung für komplexe Datentypen – alles mit einer vertrauten, intuitiven Decorator-Syntax.

Bereit, mehr zu entdecken? Probieren Sie Native Arrow UDFs noch heute auf Databricks als Teil von Databricks Runtime 18.0 aus. Um zu beginnen, ersetzen Sie einfach Ihre bestehenden Pandas UDFs durch Arrow UDFs. In den meisten Fällen sind nur wenige Zeilen Codeänderung erforderlich, um sofortige Leistungssteigerungen zu erzielen. Weitere Informationen finden Sie in der Arrow UDF-Dokumentation und der Arrow UDTF-Dokumentation für die vollständige API-Referenz und zusätzliche Beispiele.

(Dieser Blogbeitrag wurde mit KI-gestützten Tools übersetzt.) Originalbeitrag

Erhalten Sie die neuesten Beiträge in Ihrem Posteingang

Abonnieren Sie unseren Blog und erhalten Sie die neuesten Beiträge direkt in Ihren Posteingang.