Revenir au contenu principal
Open Source

Présentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus léger pour les UDF Pandas

Définissez des UDFs plus performantes en toute simplicité.

par Ruifeng Zheng et Yicong Huang

  • Nous introduisons les UDFs Arrow natives, qui opèrent directement sur les données Arrow, éliminant la surcharge de conversion Pandas/Arrow dans les UDFs Pandas pour une exécution plus rapide et une utilisation réduite de la mémoire.
  • Nous décrivons également les types UDF Arrow pour les cas d'utilisation scalaires et d'agrégation, et les UDTF Arrow pour les transformations table-en, table-out, avec des exemples de code en Python et SQL.
  • Les benchmarks montrent que les UDFs Arrow sont environ 10% plus rapides et utilisent environ 40% moins de mémoire que les UDFs Pandas, avec un meilleur support pour les types de données complexes.

Introduction

Les fonctions définies par l'utilisateur (UDF) Python sont un mécanisme d'extensibilité essentiel, mais ont traditionnellement souffert d'une surcharge élevée due à l'exécution ligne par ligne. Dans Apache Spark™, les UDF Pandas ont résolu une partie de ce problème en introduisant la sérialisation basée sur Arrow et le traitement par lots, améliorant considérablement le débit par rapport aux UDF Python scalaires.

Cependant, les UDF Pandas présentent encore des limitations fondamentales :

  • La conversion de données Pandas/Arrow introduit des copies de données supplémentaires. Les approches sans copie ne sont possibles que dans certains cas spécifiques. Par exemple, les colonnes avec des valeurs NULL déclencheront des copies profondes.
  • Les types de données complexes ne sont pas bien pris en charge. Par exemple, les instances imbriquées de StructType ne sont pas prises en charge pour le type de sortie avec des cas d'utilisation d'agrégation.
Flux de données de l'exécution des UDF Pandas dans Apache Spark

En supprimant la conversion de données Pandas/Arrow, les UDF Arrow s'exécutent plus rapidement que les UDF Pandas, consomment moins de mémoire et offrent une meilleure prise en charge des types de données.

UDF Arrow natives

Nous sommes ravis de présenter les UDF Arrow natives à partir de Databricks Runtime 18.0 (notes de mise à jour), un bond en avant passionnant pour l'exécution performante des UDF.

Les UDF Arrow natives opèrent directement sur les données Arrow sans convertir les entrées en objets Pandas ou NumPy. Cela préserve la disposition en colonnes de bout en bout, évite les copies de données inutiles et permet aux UDF d'utiliser le traitement vectorisé en tirant parti du modèle de calcul et de mémoire natif d'Arrow.

Pour définir une UDF Arrow, les utilisateurs peuvent utiliser un nouveau décorateur Python @arrow_udf, avec le type de retour spécifié et le type d'évaluation facultatif. Par exemple :

Les utilisateurs peuvent également la définir avec le décorateur existant @udf avec des indications de type complètes. Par exemple :

Remarque : La définition de la fonction doit inclure des indications de type pour tous les arguments et la valeur de retour.
Cette conception s'aligne sur les interfaces des UDF Python scalaires, offrant une expérience cohérente et intuitive aux utilisateurs déjà familiarisés avec les UDF Python scalaires.

Ce qui suit montre comment utiliser l'UDF Arrow :

Utilisation Python :

Utilisation SQL :

Nous prenons en charge des variantes d'interfaces UDF Arrow. Y compris les fonctions scalaires, les fonctions d'agrégation et les fonctions de table. Dans l'API de data frame, nous fournissons également mapInArrow et applyInArrow pour utiliser les UDF Arrow. Nous allons ensuite les présenter une par une.

Fonctions scalaires Arrow

Les fonctions scalaires Arrow effectuent des transformations ligne par ligne. Ce sont l'équivalent Arrow des UDF Pandas scalaires et peuvent être utilisées partout où une expression de colonne est attendue, comme df.select() ou df.withColumn(). Trois modes d'entrée sont pris en charge : direct, itérateur et itérateur de plusieurs tableaux. Les variantes itérateurs sont utiles lorsque l'UDF nécessite une initialisation coûteuse unique (par exemple, le chargement d'un modèle ou la compilation d'un modèle regex), car le coût de configuration est amorti sur tous les lots. Dans tous les cas, le nombre de lignes de sortie doit correspondre au nombre de lignes d'entrée.

  • Tableaux vers Tableau : réception d'un ou plusieurs pyarrow.Array et retour d'un pyarrow.Array. Le tableau d'entrée et de sortie doit avoir le même nombre de valeurs.
  • Itérateur de tableaux vers Itérateur de tableaux : réception d'un itérateur de pyarrow.Array et retour d'un itérateur de pyarrow.Array. Ce type est utile lorsque l'exécution de l'UDF nécessite une initialisation coûteuse.
  • Itérateur de tableaux multiples vers Itérateur de tableaux : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'un itérateur de pyarrow.Array.

Fonctions d'agrégation Arrow

Les fonctions d'agrégation Arrow prennent une ou plusieurs entrées pyarrow.Array et retournent une valeur scalaire, réduisant un groupe de lignes en un seul résultat. Ce sont l'équivalent Arrow des UDF Pandas agrégées groupées et sont utilisées avec groupBy().agg() ou des opérations Window. Similaires aux fonctions scalaires, les fonctions d'agrégation prennent également en charge trois modes d'entrée.

Tableaux vers Scalaire : réception de pyarrow.Array et retour d'une valeur scalaire.

  • Itérateur de tableaux vers Scalaire : réception d'un itérateur de pyarrow.Array et retour d'une valeur scalaire. Ceci est utile pour traiter de grands volumes de données dans des opérations de style d'agrégation.

Itérateur de tableaux multiples vers Scalaire : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'une valeur scalaire. Des agrégations plus complexes peuvent être définies.

Fonctions de table Arrow

Les fonctions de table Arrow, également connues sous le nom de fonctions de table définies par l'utilisateur (UDTF) Arrow, acceptent un pyarrow.RecordBatch ou plusieurs pa.Array en entrée et produisent un pyarrow.Table en sortie. Cela représente le modèle prédominant pour les transformations table-en, table-out implémentées en Python utilisant l'exécution en colonnes. Les UDTF Arrow ont la capacité de :

  • Retourner plusieurs colonnes
  • Produire zéro, une ou plusieurs lignes
  • Exécuter des transformations de table vectorisées en utilisant les noyaux de calcul Arrow

Par conséquent, elles sont idéalement adaptées aux opérations telles que le filtrage, l'expansion de lignes, la restructuration de données et la génération de colonnes dérivées.

L'interface arrow_udtf est conçue pour la simplicité, utilisant une syntaxe de décorateur où vous définissez le type de retour à l'aide d'une chaîne formatée en DDL. Dans cette configuration, la méthode eval prend des objets PyArrow en entrée et est censée générer des PyArrow Tables ou RecordBatches. L'interface prend en charge deux modes d'entrée. Lors du traitement des arguments de table, la méthode eval reçoit un objet pa.RecordBatch qui encapsule toutes les colonnes de la table d'entrée :

Pour les arguments scalaires, la méthode reçoit des objets pa.Array, un pour chaque entrée scalaire :

Voici un autre exemple :

Cette UDTF peut fonctionner de deux manières distinctes :

Utilisation Python :

Utilisation SQL :

Prise en charge de DataFrame mapInArrow et applyInArrow

En plus des fonctions définies par l'utilisateur (UDF) et des fonctions de table définies par l'utilisateur (UDTF), PySpark fournit des API de fonctions Arrow qui facilitent l'application directe de fonctions Python natives aux données Arrow au niveau du DataFrame. Ces API fonctionnent de manière analogue à leurs homologues Pandas (mapInPandas, applyInPandas) mais utilisent pyarrow.RecordBatch et pyarrow.Table au lieu des DataFrames Pandas, évitant ainsi les frais de conversion entre les formats Pandas et Arrow.

  • Map. DataFrame.mapInArrow transforme un itérateur de pyarrow.RecordBatch en un autre itérateur de pyarrow.RecordBatch, permettant des opérations au niveau des lignes telles que le filtrage, la transformation ou l'expansion.
  • Grouped Map. groupBy().applyInArrow() applique une fonction spécifiée à chaque groupe, acceptant et retournant un pyarrow.Table. Cette fonctionnalité est utile pour les transformations par groupe, telles que la normalisation des données.
  • Co-grouped Map. cogroup().applyInArrow() permet le co-groupement de deux DataFrames sur la base d'une clé partagée, puis l'application d'une fonction à chaque co-groupe. La fonction reçoit deux entrées pyarrow.Table et est censée retourner une seule pyarrow.Table.

Performance

En supprimant la coûteuse conversion de données Pandas/Arrow, les UDF Arrow s'exécutent généralement plus rapidement que les UDF Pandas, avec une utilisation de mémoire réduite. Comparons les deux UDF simples :

L'UDF Arrow est environ 10 % plus rapide que l'UDF Pandas, et le profileur de mémoire montre qu'environ 40 % de mémoire est économisée lors de l'exécution.

Conclusion

Databricks Runtime 18.0 introduit les UDF Arrow natives, offrant une alternative plus rapide et plus légère aux UDF Pandas pour une exécution performante des UDF Python dans PySpark. En opérant directement sur les données Arrow et en éliminant les frais de conversion Pandas/Arrow, les UDF Arrow offrent une exécution environ 10 % plus rapide, une utilisation de mémoire environ 40 % inférieure et un meilleur support pour les types de données complexes, le tout avec une syntaxe de décorateur familière et intuitive.

Prêt à en découvrir davantage ? Essayez dès aujourd'hui les UDF Arrow natives sur Databricks dans le cadre de Databricks Runtime 18.0. Pour commencer, remplacez simplement vos UDF Pandas existantes par des UDF Arrow. Dans la plupart des cas, il suffit de quelques lignes de modification pour obtenir des gains de performance immédiats. Consultez la documentation des UDF Arrow et la documentation des UDTF pour la référence API complète et des exemples supplémentaires.

(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original

Recevez les derniers articles dans votre boîte mail

Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])It\u00e9rateur de tableaux multiples vers It\u00e9rateur de tableaux : r\u00e9ception d'un it\u00e9rateur d'un tuple de plusieurs pyarrow.Array et retour d'un it\u00e9rateur de 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)Fonctions d'agr\u00e9gation ArrowLes fonctions d'agr\u00e9gation Arrow prennent une ou plusieurs entr\u00e9es pyarrow.Array et retournent une valeur scalaire, r\u00e9duisant un groupe de lignes en un seul r\u00e9sultat. Ce sont l'\u00e9quivalent Arrow des UDF Pandas agr\u00e9g\u00e9es group\u00e9es et sont utilis\u00e9es avec groupBy().agg() ou des op\u00e9rations Window. Similaires aux fonctions scalaires, les fonctions d'agr\u00e9gation prennent \u00e9galement en charge trois modes d'entr\u00e9e. Tableaux vers Scalaire : r\u00e9ception de pyarrow.Array et retour d'une valeur scalaire. 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)It\u00e9rateur de tableaux vers Scalaire : r\u00e9ception d'un it\u00e9rateur de pyarrow.Array et retour d'une valeur scalaire. Ceci est utile pour traiter de grands volumes de donn\u00e9es dans des op\u00e9rations de style d'agr\u00e9gation.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.0It\u00e9rateur de tableaux multiples vers Scalaire : r\u00e9ception d'un it\u00e9rateur d'un tuple de plusieurs pyarrow.Array et retour d'une valeur scalaire. Des agr\u00e9gations plus complexes peuvent \u00eatre d\u00e9finies.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.0Fonctions de table ArrowLes fonctions de table Arrow, \u00e9galement connues sous le nom de fonctions de table d\u00e9finies par l'utilisateur (UDTF) Arrow, acceptent un pyarrow.RecordBatch ou plusieurs pa.Array en entr\u00e9e et produisent un pyarrow.Table en sortie. Cela repr\u00e9sente le mod\u00e8le pr\u00e9dominant pour les transformations table-en, table-out impl\u00e9ment\u00e9es en Python utilisant l'ex\u00e9cution en colonnes. Les UDTF Arrow ont la capacit\u00e9 de :Retourner plusieurs colonnesProduire z\u00e9ro, une ou plusieurs lignesEx\u00e9cuter des transformations de table vectoris\u00e9es en utilisant les noyaux de calcul ArrowPar cons\u00e9quent, elles sont id\u00e9alement adapt\u00e9es aux op\u00e9rations telles que le filtrage, l'expansion de lignes, la restructuration de donn\u00e9es et la g\u00e9n\u00e9ration de colonnes d\u00e9riv\u00e9es.L'interface arrow_udtf est con\u00e7ue pour la simplicit\u00e9, utilisant une syntaxe de d\u00e9corateur o\u00f9 vous d\u00e9finissez le type de retour \u00e0 l'aide d'une cha\u00eene format\u00e9e en DDL. Dans cette configuration, la m\u00e9thode eval prend des objets PyArrow en entr\u00e9e et est cens\u00e9e g\u00e9n\u00e9rer des PyArrow Tables ou RecordBatches. L'interface prend en charge deux modes d'entr\u00e9e. Lors du traitement des arguments de table, la m\u00e9thode eval re\u00e7oit un objet pa.RecordBatch qui encapsule toutes les colonnes de la table d'entr\u00e9e :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 resultPour les arguments scalaires, la m\u00e9thode re\u00e7oit des objets pa.Array, un pour chaque entr\u00e9e scalaire :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 resultVoici un autre exemple :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 })Cette UDTF peut fonctionner de deux mani\u00e8res distinctes :Utilisation 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# +----------+-------+--------+Utilisation 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# +----------+-------+--------+Prise en charge de DataFrame mapInArrow et applyInArrowEn plus des fonctions d\u00e9finies par l'utilisateur (UDF) et des fonctions de table d\u00e9finies par l'utilisateur (UDTF), PySpark fournit des API de fonctions Arrow qui facilitent l'application directe de fonctions Python natives aux donn\u00e9es Arrow au niveau du DataFrame. Ces API fonctionnent de mani\u00e8re analogue \u00e0 leurs homologues Pandas (mapInPandas, applyInPandas) mais utilisent pyarrow.RecordBatch et pyarrow.Table au lieu des DataFrames Pandas, \u00e9vitant ainsi les frais de conversion entre les formats Pandas et Arrow.Map. DataFrame.mapInArrow transforme un it\u00e9rateur de pyarrow.RecordBatch en un autre it\u00e9rateur de pyarrow.RecordBatch, permettant des op\u00e9rations au niveau des lignes telles que le filtrage, la transformation ou l'expansion.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() applique une fonction sp\u00e9cifi\u00e9e \u00e0 chaque groupe, acceptant et retournant un pyarrow.Table. Cette fonctionnalit\u00e9 est utile pour les transformations par groupe, telles que la normalisation des donn\u00e9es.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() permet le co-groupement de deux DataFrames sur la base d'une cl\u00e9 partag\u00e9e, puis l'application d'une fonction \u00e0 chaque co-groupe. La fonction re\u00e7oit deux entr\u00e9es pyarrow.Table et est cens\u00e9e retourner une seule 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# +----+-----+PerformanceEn supprimant la co\u00fbteuse conversion de donn\u00e9es Pandas/Arrow, les UDF Arrow s'ex\u00e9cutent g\u00e9n\u00e9ralement plus rapidement que les UDF Pandas, avec une utilisation de m\u00e9moire r\u00e9duite. Comparons les deux UDF simples :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 est environ 10 % plus rapide que l'UDF Pandas, et le profileur de m\u00e9moire montre qu'environ 40 % de m\u00e9moire est \u00e9conomis\u00e9e lors de l'ex\u00e9cution.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)ConclusionDatabricks Runtime 18.0 introduit les UDF Arrow natives, offrant une alternative plus rapide et plus l\u00e9g\u00e8re aux UDF Pandas pour une ex\u00e9cution performante des UDF Python dans PySpark. En op\u00e9rant directement sur les donn\u00e9es Arrow et en \u00e9liminant les frais de conversion Pandas/Arrow, les UDF Arrow offrent une ex\u00e9cution environ 10 % plus rapide, une utilisation de m\u00e9moire environ 40 % inf\u00e9rieure et un meilleur support pour les types de donn\u00e9es complexes, le tout avec une syntaxe de d\u00e9corateur famili\u00e8re et intuitive.Pr\u00eat \u00e0 en d\u00e9couvrir davantage ? Essayez d\u00e8s aujourd'hui les UDF Arrow natives sur Databricks dans le cadre de Databricks Runtime 18.0. Pour commencer, remplacez simplement vos UDF Pandas existantes par des UDF Arrow. Dans la plupart des cas, il suffit de quelques lignes de modification pour obtenir des gains de performance imm\u00e9diats. Consultez la documentation des UDF Arrow et la documentation des UDTF pour la r\u00e9f\u00e9rence API compl\u00e8te et des exemples suppl\u00e9mentaires.(Cet article de blog a \u00e9t\u00e9 traduit \u00e0 l'aide d'outils bas\u00e9s sur l'intelligence artificielle) Article original", "description": "D\u00e9couvrez comment \u00e9crire des UDFs plus performantes avec le support natif de pyarrow.", "headline": "Pr\u00e9sentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus l\u00e9ger pour les UDF Pandas", "mainEntityOfPage": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/fr/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"}], "mentions": [{"@type": "BreadcrumbList", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "item": "https://www.databricks.com/fr/blog", "name": "Tous les blogs", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "item": "https://www.databricks.com/fr/blog/category/engineering", "name": "Ing\u00e9nierie", "position": 2}]}], "articleSection": "Connexion\nOpen Source", "author": [{"@type": "Person", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/fr/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/fr/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"}]
Revenir au contenu principal
Open Source

Présentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus léger pour les UDF Pandas

Définissez des UDFs plus performantes en toute simplicité.

par Ruifeng Zheng et Yicong Huang

  • Nous introduisons les UDFs Arrow natives, qui opèrent directement sur les données Arrow, éliminant la surcharge de conversion Pandas/Arrow dans les UDFs Pandas pour une exécution plus rapide et une utilisation réduite de la mémoire.
  • Nous décrivons également les types UDF Arrow pour les cas d'utilisation scalaires et d'agrégation, et les UDTF Arrow pour les transformations table-en, table-out, avec des exemples de code en Python et SQL.
  • Les benchmarks montrent que les UDFs Arrow sont environ 10% plus rapides et utilisent environ 40% moins de mémoire que les UDFs Pandas, avec un meilleur support pour les types de données complexes.

Introduction

Les fonctions définies par l'utilisateur (UDF) Python sont un mécanisme d'extensibilité essentiel, mais ont traditionnellement souffert d'une surcharge élevée due à l'exécution ligne par ligne. Dans Apache Spark™, les UDF Pandas ont résolu une partie de ce problème en introduisant la sérialisation basée sur Arrow et le traitement par lots, améliorant considérablement le débit par rapport aux UDF Python scalaires.

Cependant, les UDF Pandas présentent encore des limitations fondamentales :

  • La conversion de données Pandas/Arrow introduit des copies de données supplémentaires. Les approches sans copie ne sont possibles que dans certains cas spécifiques. Par exemple, les colonnes avec des valeurs NULL déclencheront des copies profondes.
  • Les types de données complexes ne sont pas bien pris en charge. Par exemple, les instances imbriquées de StructType ne sont pas prises en charge pour le type de sortie avec des cas d'utilisation d'agrégation.
Flux de données de l'exécution des UDF Pandas dans Apache Spark

En supprimant la conversion de données Pandas/Arrow, les UDF Arrow s'exécutent plus rapidement que les UDF Pandas, consomment moins de mémoire et offrent une meilleure prise en charge des types de données.

UDF Arrow natives

Nous sommes ravis de présenter les UDF Arrow natives à partir de Databricks Runtime 18.0 (notes de mise à jour), un bond en avant passionnant pour l'exécution performante des UDF.

Les UDF Arrow natives opèrent directement sur les données Arrow sans convertir les entrées en objets Pandas ou NumPy. Cela préserve la disposition en colonnes de bout en bout, évite les copies de données inutiles et permet aux UDF d'utiliser le traitement vectorisé en tirant parti du modèle de calcul et de mémoire natif d'Arrow.

Pour définir une UDF Arrow, les utilisateurs peuvent utiliser un nouveau décorateur Python @arrow_udf, avec le type de retour spécifié et le type d'évaluation facultatif. Par exemple :

Les utilisateurs peuvent également la définir avec le décorateur existant @udf avec des indications de type complètes. Par exemple :

Remarque : La définition de la fonction doit inclure des indications de type pour tous les arguments et la valeur de retour.
Cette conception s'aligne sur les interfaces des UDF Python scalaires, offrant une expérience cohérente et intuitive aux utilisateurs déjà familiarisés avec les UDF Python scalaires.

Ce qui suit montre comment utiliser l'UDF Arrow :

Utilisation Python :

Utilisation SQL :

Nous prenons en charge des variantes d'interfaces UDF Arrow. Y compris les fonctions scalaires, les fonctions d'agrégation et les fonctions de table. Dans l'API de data frame, nous fournissons également mapInArrow et applyInArrow pour utiliser les UDF Arrow. Nous allons ensuite les présenter une par une.

Fonctions scalaires Arrow

Les fonctions scalaires Arrow effectuent des transformations ligne par ligne. Ce sont l'équivalent Arrow des UDF Pandas scalaires et peuvent être utilisées partout où une expression de colonne est attendue, comme df.select() ou df.withColumn(). Trois modes d'entrée sont pris en charge : direct, itérateur et itérateur de plusieurs tableaux. Les variantes itérateurs sont utiles lorsque l'UDF nécessite une initialisation coûteuse unique (par exemple, le chargement d'un modèle ou la compilation d'un modèle regex), car le coût de configuration est amorti sur tous les lots. Dans tous les cas, le nombre de lignes de sortie doit correspondre au nombre de lignes d'entrée.

  • Tableaux vers Tableau : réception d'un ou plusieurs pyarrow.Array et retour d'un pyarrow.Array. Le tableau d'entrée et de sortie doit avoir le même nombre de valeurs.
  • Itérateur de tableaux vers Itérateur de tableaux : réception d'un itérateur de pyarrow.Array et retour d'un itérateur de pyarrow.Array. Ce type est utile lorsque l'exécution de l'UDF nécessite une initialisation coûteuse.
  • Itérateur de tableaux multiples vers Itérateur de tableaux : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'un itérateur de pyarrow.Array.

Fonctions d'agrégation Arrow

Les fonctions d'agrégation Arrow prennent une ou plusieurs entrées pyarrow.Array et retournent une valeur scalaire, réduisant un groupe de lignes en un seul résultat. Ce sont l'équivalent Arrow des UDF Pandas agrégées groupées et sont utilisées avec groupBy().agg() ou des opérations Window. Similaires aux fonctions scalaires, les fonctions d'agrégation prennent également en charge trois modes d'entrée.

Tableaux vers Scalaire : réception de pyarrow.Array et retour d'une valeur scalaire.

  • Itérateur de tableaux vers Scalaire : réception d'un itérateur de pyarrow.Array et retour d'une valeur scalaire. Ceci est utile pour traiter de grands volumes de données dans des opérations de style d'agrégation.

Itérateur de tableaux multiples vers Scalaire : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'une valeur scalaire. Des agrégations plus complexes peuvent être définies.

Fonctions de table Arrow

Les fonctions de table Arrow, également connues sous le nom de fonctions de table définies par l'utilisateur (UDTF) Arrow, acceptent un pyarrow.RecordBatch ou plusieurs pa.Array en entrée et produisent un pyarrow.Table en sortie. Cela représente le modèle prédominant pour les transformations table-en, table-out implémentées en Python utilisant l'exécution en colonnes. Les UDTF Arrow ont la capacité de :

  • Retourner plusieurs colonnes
  • Produire zéro, une ou plusieurs lignes
  • Exécuter des transformations de table vectorisées en utilisant les noyaux de calcul Arrow

Par conséquent, elles sont idéalement adaptées aux opérations telles que le filtrage, l'expansion de lignes, la restructuration de données et la génération de colonnes dérivées.

L'interface arrow_udtf est conçue pour la simplicité, utilisant une syntaxe de décorateur où vous définissez le type de retour à l'aide d'une chaîne formatée en DDL. Dans cette configuration, la méthode eval prend des objets PyArrow en entrée et est censée générer des PyArrow Tables ou RecordBatches. L'interface prend en charge deux modes d'entrée. Lors du traitement des arguments de table, la méthode eval reçoit un objet pa.RecordBatch qui encapsule toutes les colonnes de la table d'entrée :

Pour les arguments scalaires, la méthode reçoit des objets pa.Array, un pour chaque entrée scalaire :

Voici un autre exemple :

Cette UDTF peut fonctionner de deux manières distinctes :

Utilisation Python :

Utilisation SQL :

Prise en charge de DataFrame mapInArrow et applyInArrow

En plus des fonctions définies par l'utilisateur (UDF) et des fonctions de table définies par l'utilisateur (UDTF), PySpark fournit des API de fonctions Arrow qui facilitent l'application directe de fonctions Python natives aux données Arrow au niveau du DataFrame. Ces API fonctionnent de manière analogue à leurs homologues Pandas (mapInPandas, applyInPandas) mais utilisent pyarrow.RecordBatch et pyarrow.Table au lieu des DataFrames Pandas, évitant ainsi les frais de conversion entre les formats Pandas et Arrow.

  • Map. DataFrame.mapInArrow transforme un itérateur de pyarrow.RecordBatch en un autre itérateur de pyarrow.RecordBatch, permettant des opérations au niveau des lignes telles que le filtrage, la transformation ou l'expansion.
  • Grouped Map. groupBy().applyInArrow() applique une fonction spécifiée à chaque groupe, acceptant et retournant un pyarrow.Table. Cette fonctionnalité est utile pour les transformations par groupe, telles que la normalisation des données.
  • Co-grouped Map. cogroup().applyInArrow() permet le co-groupement de deux DataFrames sur la base d'une clé partagée, puis l'application d'une fonction à chaque co-groupe. La fonction reçoit deux entrées pyarrow.Table et est censée retourner une seule pyarrow.Table.

Performance

En supprimant la coûteuse conversion de données Pandas/Arrow, les UDF Arrow s'exécutent généralement plus rapidement que les UDF Pandas, avec une utilisation de mémoire réduite. Comparons les deux UDF simples :

L'UDF Arrow est environ 10 % plus rapide que l'UDF Pandas, et le profileur de mémoire montre qu'environ 40 % de mémoire est économisée lors de l'exécution.

Conclusion

Databricks Runtime 18.0 introduit les UDF Arrow natives, offrant une alternative plus rapide et plus légère aux UDF Pandas pour une exécution performante des UDF Python dans PySpark. En opérant directement sur les données Arrow et en éliminant les frais de conversion Pandas/Arrow, les UDF Arrow offrent une exécution environ 10 % plus rapide, une utilisation de mémoire environ 40 % inférieure et un meilleur support pour les types de données complexes, le tout avec une syntaxe de décorateur familière et intuitive.

Prêt à en découvrir davantage ? Essayez dès aujourd'hui les UDF Arrow natives sur Databricks dans le cadre de Databricks Runtime 18.0. Pour commencer, remplacez simplement vos UDF Pandas existantes par des UDF Arrow. Dans la plupart des cas, il suffit de quelques lignes de modification pour obtenir des gains de performance immédiats. Consultez la documentation des UDF Arrow et la documentation des UDTF pour la référence API complète et des exemples supplémentaires.

(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original

Recevez les derniers articles dans votre boîte mail

Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])It\u00e9rateur de tableaux multiples vers It\u00e9rateur de tableaux : r\u00e9ception d'un it\u00e9rateur d'un tuple de plusieurs pyarrow.Array et retour d'un it\u00e9rateur de 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)Fonctions d'agr\u00e9gation ArrowLes fonctions d'agr\u00e9gation Arrow prennent une ou plusieurs entr\u00e9es pyarrow.Array et retournent une valeur scalaire, r\u00e9duisant un groupe de lignes en un seul r\u00e9sultat. Ce sont l'\u00e9quivalent Arrow des UDF Pandas agr\u00e9g\u00e9es group\u00e9es et sont utilis\u00e9es avec groupBy().agg() ou des op\u00e9rations Window. Similaires aux fonctions scalaires, les fonctions d'agr\u00e9gation prennent \u00e9galement en charge trois modes d'entr\u00e9e. Tableaux vers Scalaire : r\u00e9ception de pyarrow.Array et retour d'une valeur scalaire. 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)It\u00e9rateur de tableaux vers Scalaire : r\u00e9ception d'un it\u00e9rateur de pyarrow.Array et retour d'une valeur scalaire. Ceci est utile pour traiter de grands volumes de donn\u00e9es dans des op\u00e9rations de style d'agr\u00e9gation.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.0It\u00e9rateur de tableaux multiples vers Scalaire : r\u00e9ception d'un it\u00e9rateur d'un tuple de plusieurs pyarrow.Array et retour d'une valeur scalaire. Des agr\u00e9gations plus complexes peuvent \u00eatre d\u00e9finies.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.0Fonctions de table ArrowLes fonctions de table Arrow, \u00e9galement connues sous le nom de fonctions de table d\u00e9finies par l'utilisateur (UDTF) Arrow, acceptent un pyarrow.RecordBatch ou plusieurs pa.Array en entr\u00e9e et produisent un pyarrow.Table en sortie. Cela repr\u00e9sente le mod\u00e8le pr\u00e9dominant pour les transformations table-en, table-out impl\u00e9ment\u00e9es en Python utilisant l'ex\u00e9cution en colonnes. Les UDTF Arrow ont la capacit\u00e9 de :Retourner plusieurs colonnesProduire z\u00e9ro, une ou plusieurs lignesEx\u00e9cuter des transformations de table vectoris\u00e9es en utilisant les noyaux de calcul ArrowPar cons\u00e9quent, elles sont id\u00e9alement adapt\u00e9es aux op\u00e9rations telles que le filtrage, l'expansion de lignes, la restructuration de donn\u00e9es et la g\u00e9n\u00e9ration de colonnes d\u00e9riv\u00e9es.L'interface arrow_udtf est con\u00e7ue pour la simplicit\u00e9, utilisant une syntaxe de d\u00e9corateur o\u00f9 vous d\u00e9finissez le type de retour \u00e0 l'aide d'une cha\u00eene format\u00e9e en DDL. Dans cette configuration, la m\u00e9thode eval prend des objets PyArrow en entr\u00e9e et est cens\u00e9e g\u00e9n\u00e9rer des PyArrow Tables ou RecordBatches. L'interface prend en charge deux modes d'entr\u00e9e. Lors du traitement des arguments de table, la m\u00e9thode eval re\u00e7oit un objet pa.RecordBatch qui encapsule toutes les colonnes de la table d'entr\u00e9e :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 resultPour les arguments scalaires, la m\u00e9thode re\u00e7oit des objets pa.Array, un pour chaque entr\u00e9e scalaire :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 resultVoici un autre exemple :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 })Cette UDTF peut fonctionner de deux mani\u00e8res distinctes :Utilisation 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# +----------+-------+--------+Utilisation 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# +----------+-------+--------+Prise en charge de DataFrame mapInArrow et applyInArrowEn plus des fonctions d\u00e9finies par l'utilisateur (UDF) et des fonctions de table d\u00e9finies par l'utilisateur (UDTF), PySpark fournit des API de fonctions Arrow qui facilitent l'application directe de fonctions Python natives aux donn\u00e9es Arrow au niveau du DataFrame. Ces API fonctionnent de mani\u00e8re analogue \u00e0 leurs homologues Pandas (mapInPandas, applyInPandas) mais utilisent pyarrow.RecordBatch et pyarrow.Table au lieu des DataFrames Pandas, \u00e9vitant ainsi les frais de conversion entre les formats Pandas et Arrow.Map. DataFrame.mapInArrow transforme un it\u00e9rateur de pyarrow.RecordBatch en un autre it\u00e9rateur de pyarrow.RecordBatch, permettant des op\u00e9rations au niveau des lignes telles que le filtrage, la transformation ou l'expansion.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() applique une fonction sp\u00e9cifi\u00e9e \u00e0 chaque groupe, acceptant et retournant un pyarrow.Table. Cette fonctionnalit\u00e9 est utile pour les transformations par groupe, telles que la normalisation des donn\u00e9es.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() permet le co-groupement de deux DataFrames sur la base d'une cl\u00e9 partag\u00e9e, puis l'application d'une fonction \u00e0 chaque co-groupe. La fonction re\u00e7oit deux entr\u00e9es pyarrow.Table et est cens\u00e9e retourner une seule 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# +----+-----+PerformanceEn supprimant la co\u00fbteuse conversion de donn\u00e9es Pandas/Arrow, les UDF Arrow s'ex\u00e9cutent g\u00e9n\u00e9ralement plus rapidement que les UDF Pandas, avec une utilisation de m\u00e9moire r\u00e9duite. Comparons les deux UDF simples :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 est environ 10 % plus rapide que l'UDF Pandas, et le profileur de m\u00e9moire montre qu'environ 40 % de m\u00e9moire est \u00e9conomis\u00e9e lors de l'ex\u00e9cution.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)ConclusionDatabricks Runtime 18.0 introduit les UDF Arrow natives, offrant une alternative plus rapide et plus l\u00e9g\u00e8re aux UDF Pandas pour une ex\u00e9cution performante des UDF Python dans PySpark. En op\u00e9rant directement sur les donn\u00e9es Arrow et en \u00e9liminant les frais de conversion Pandas/Arrow, les UDF Arrow offrent une ex\u00e9cution environ 10 % plus rapide, une utilisation de m\u00e9moire environ 40 % inf\u00e9rieure et un meilleur support pour les types de donn\u00e9es complexes, le tout avec une syntaxe de d\u00e9corateur famili\u00e8re et intuitive.Pr\u00eat \u00e0 en d\u00e9couvrir davantage ? Essayez d\u00e8s aujourd'hui les UDF Arrow natives sur Databricks dans le cadre de Databricks Runtime 18.0. Pour commencer, remplacez simplement vos UDF Pandas existantes par des UDF Arrow. Dans la plupart des cas, il suffit de quelques lignes de modification pour obtenir des gains de performance imm\u00e9diats. Consultez la documentation des UDF Arrow et la documentation des UDTF pour la r\u00e9f\u00e9rence API compl\u00e8te et des exemples suppl\u00e9mentaires.(Cet article de blog a \u00e9t\u00e9 traduit \u00e0 l'aide d'outils bas\u00e9s sur l'intelligence artificielle) Article original", "headline": "Pr\u00e9sentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus l\u00e9ger pour les UDF Pandas", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/fr/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"}], "mainEntityOfPage": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs", "name": "Pr\u00e9sentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus l\u00e9ger pour les UDF Pandas", "articleSection": "Open Source", "author": [{"@type": "Person", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/fr/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/fr/blog/author/yicong-huang", "name": "Yicong Huang"}], "mentions": [{"@type": "BreadcrumbList", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Tous les blogs", "item": "https://www.databricks.com/fr/blog", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/fr/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Ing\u00e9nierie", "item": "https://www.databricks.com/fr/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": "run-time system", "@id": "https://entity.schemaapp.com/DatabricksInc/Thing_run-timesystem_0cc5c91b63283929ce49e5094aacd5b5f16262aa986f2a16b874af0902ba2a6a", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/06mj_w", "http://www.wikidata.org/entity/Q1004415", "https://en.wikipedia.org/wiki/Runtime_system"]}, {"name": "Parti Radical", "@id": "https://entity.schemaapp.com/DatabricksInc/Organization_partiradical_4b392e39a2cb87f33d453760f3ba0ab7601f333473042c0b9532c51481467b0d", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/022f4r", "http://www.wikidata.org/entity/Q1542710", "https://en.wikipedia.org/wiki/Radical_Party_(France)"]}, {"name": "Union pour la democratie francaise", "@id": "https://entity.schemaapp.com/DatabricksInc/Organization_unionpourlademocratiefrancaise_44e1392ba31dbc06fc40a3a76168c6599de71736cc4a2cd77647b9afe400f03a", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/01vv7g", "http://www.wikidata.org/entity/Q827415", "https://en.wikipedia.org/wiki/Union_for_French_Democracy"]}, {"name": "United States", "@id": "https://entity.schemaapp.com/DatabricksInc/Country_unitedstates_35761fb7e035aa623e469abc3c609c0eebb9ce047778da1e53d49bf8709aa9c5", "@type": "Thing", "@context": {"@vocab": "http://schema.org/"}, "sameAs": ["http://g.co/kg/m/09c7w0", "http://www.wikidata.org/entity/Q30", "https://en.wikipedia.org/wiki/United_States"]}]}, {"@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"}]
Revenir au contenu principal
Open Source

Présentation des UDF Arrow dans PySpark : un remplacement plus rapide et plus léger pour les UDF Pandas

Définissez des UDFs plus performantes en toute simplicité.

par Ruifeng Zheng et Yicong Huang

  • Nous introduisons les UDFs Arrow natives, qui opèrent directement sur les données Arrow, éliminant la surcharge de conversion Pandas/Arrow dans les UDFs Pandas pour une exécution plus rapide et une utilisation réduite de la mémoire.
  • Nous décrivons également les types UDF Arrow pour les cas d'utilisation scalaires et d'agrégation, et les UDTF Arrow pour les transformations table-en, table-out, avec des exemples de code en Python et SQL.
  • Les benchmarks montrent que les UDFs Arrow sont environ 10% plus rapides et utilisent environ 40% moins de mémoire que les UDFs Pandas, avec un meilleur support pour les types de données complexes.

Introduction

Les fonctions définies par l'utilisateur (UDF) Python sont un mécanisme d'extensibilité essentiel, mais ont traditionnellement souffert d'une surcharge élevée due à l'exécution ligne par ligne. Dans Apache Spark™, les UDF Pandas ont résolu une partie de ce problème en introduisant la sérialisation basée sur Arrow et le traitement par lots, améliorant considérablement le débit par rapport aux UDF Python scalaires.

Cependant, les UDF Pandas présentent encore des limitations fondamentales :

  • La conversion de données Pandas/Arrow introduit des copies de données supplémentaires. Les approches sans copie ne sont possibles que dans certains cas spécifiques. Par exemple, les colonnes avec des valeurs NULL déclencheront des copies profondes.
  • Les types de données complexes ne sont pas bien pris en charge. Par exemple, les instances imbriquées de StructType ne sont pas prises en charge pour le type de sortie avec des cas d'utilisation d'agrégation.
Flux de données de l'exécution des UDF Pandas dans Apache Spark

En supprimant la conversion de données Pandas/Arrow, les UDF Arrow s'exécutent plus rapidement que les UDF Pandas, consomment moins de mémoire et offrent une meilleure prise en charge des types de données.

UDF Arrow natives

Nous sommes ravis de présenter les UDF Arrow natives à partir de Databricks Runtime 18.0 (notes de mise à jour), un bond en avant passionnant pour l'exécution performante des UDF.

Les UDF Arrow natives opèrent directement sur les données Arrow sans convertir les entrées en objets Pandas ou NumPy. Cela préserve la disposition en colonnes de bout en bout, évite les copies de données inutiles et permet aux UDF d'utiliser le traitement vectorisé en tirant parti du modèle de calcul et de mémoire natif d'Arrow.

Pour définir une UDF Arrow, les utilisateurs peuvent utiliser un nouveau décorateur Python @arrow_udf, avec le type de retour spécifié et le type d'évaluation facultatif. Par exemple :

Les utilisateurs peuvent également la définir avec le décorateur existant @udf avec des indications de type complètes. Par exemple :

Remarque : La définition de la fonction doit inclure des indications de type pour tous les arguments et la valeur de retour.
Cette conception s'aligne sur les interfaces des UDF Python scalaires, offrant une expérience cohérente et intuitive aux utilisateurs déjà familiarisés avec les UDF Python scalaires.

Ce qui suit montre comment utiliser l'UDF Arrow :

Utilisation Python :

Utilisation SQL :

Nous prenons en charge des variantes d'interfaces UDF Arrow. Y compris les fonctions scalaires, les fonctions d'agrégation et les fonctions de table. Dans l'API de data frame, nous fournissons également mapInArrow et applyInArrow pour utiliser les UDF Arrow. Nous allons ensuite les présenter une par une.

Fonctions scalaires Arrow

Les fonctions scalaires Arrow effectuent des transformations ligne par ligne. Ce sont l'équivalent Arrow des UDF Pandas scalaires et peuvent être utilisées partout où une expression de colonne est attendue, comme df.select() ou df.withColumn(). Trois modes d'entrée sont pris en charge : direct, itérateur et itérateur de plusieurs tableaux. Les variantes itérateurs sont utiles lorsque l'UDF nécessite une initialisation coûteuse unique (par exemple, le chargement d'un modèle ou la compilation d'un modèle regex), car le coût de configuration est amorti sur tous les lots. Dans tous les cas, le nombre de lignes de sortie doit correspondre au nombre de lignes d'entrée.

  • Tableaux vers Tableau : réception d'un ou plusieurs pyarrow.Array et retour d'un pyarrow.Array. Le tableau d'entrée et de sortie doit avoir le même nombre de valeurs.
  • Itérateur de tableaux vers Itérateur de tableaux : réception d'un itérateur de pyarrow.Array et retour d'un itérateur de pyarrow.Array. Ce type est utile lorsque l'exécution de l'UDF nécessite une initialisation coûteuse.
  • Itérateur de tableaux multiples vers Itérateur de tableaux : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'un itérateur de pyarrow.Array.

Fonctions d'agrégation Arrow

Les fonctions d'agrégation Arrow prennent une ou plusieurs entrées pyarrow.Array et retournent une valeur scalaire, réduisant un groupe de lignes en un seul résultat. Ce sont l'équivalent Arrow des UDF Pandas agrégées groupées et sont utilisées avec groupBy().agg() ou des opérations Window. Similaires aux fonctions scalaires, les fonctions d'agrégation prennent également en charge trois modes d'entrée.

Tableaux vers Scalaire : réception de pyarrow.Array et retour d'une valeur scalaire.

  • Itérateur de tableaux vers Scalaire : réception d'un itérateur de pyarrow.Array et retour d'une valeur scalaire. Ceci est utile pour traiter de grands volumes de données dans des opérations de style d'agrégation.

Itérateur de tableaux multiples vers Scalaire : réception d'un itérateur d'un tuple de plusieurs pyarrow.Array et retour d'une valeur scalaire. Des agrégations plus complexes peuvent être définies.

Fonctions de table Arrow

Les fonctions de table Arrow, également connues sous le nom de fonctions de table définies par l'utilisateur (UDTF) Arrow, acceptent un pyarrow.RecordBatch ou plusieurs pa.Array en entrée et produisent un pyarrow.Table en sortie. Cela représente le modèle prédominant pour les transformations table-en, table-out implémentées en Python utilisant l'exécution en colonnes. Les UDTF Arrow ont la capacité de :

  • Retourner plusieurs colonnes
  • Produire zéro, une ou plusieurs lignes
  • Exécuter des transformations de table vectorisées en utilisant les noyaux de calcul Arrow

Par conséquent, elles sont idéalement adaptées aux opérations telles que le filtrage, l'expansion de lignes, la restructuration de données et la génération de colonnes dérivées.

L'interface arrow_udtf est conçue pour la simplicité, utilisant une syntaxe de décorateur où vous définissez le type de retour à l'aide d'une chaîne formatée en DDL. Dans cette configuration, la méthode eval prend des objets PyArrow en entrée et est censée générer des PyArrow Tables ou RecordBatches. L'interface prend en charge deux modes d'entrée. Lors du traitement des arguments de table, la méthode eval reçoit un objet pa.RecordBatch qui encapsule toutes les colonnes de la table d'entrée :

Pour les arguments scalaires, la méthode reçoit des objets pa.Array, un pour chaque entrée scalaire :

Voici un autre exemple :

Cette UDTF peut fonctionner de deux manières distinctes :

Utilisation Python :

Utilisation SQL :

Prise en charge de DataFrame mapInArrow et applyInArrow

En plus des fonctions définies par l'utilisateur (UDF) et des fonctions de table définies par l'utilisateur (UDTF), PySpark fournit des API de fonctions Arrow qui facilitent l'application directe de fonctions Python natives aux données Arrow au niveau du DataFrame. Ces API fonctionnent de manière analogue à leurs homologues Pandas (mapInPandas, applyInPandas) mais utilisent pyarrow.RecordBatch et pyarrow.Table au lieu des DataFrames Pandas, évitant ainsi les frais de conversion entre les formats Pandas et Arrow.

  • Map. DataFrame.mapInArrow transforme un itérateur de pyarrow.RecordBatch en un autre itérateur de pyarrow.RecordBatch, permettant des opérations au niveau des lignes telles que le filtrage, la transformation ou l'expansion.
  • Grouped Map. groupBy().applyInArrow() applique une fonction spécifiée à chaque groupe, acceptant et retournant un pyarrow.Table. Cette fonctionnalité est utile pour les transformations par groupe, telles que la normalisation des données.
  • Co-grouped Map. cogroup().applyInArrow() permet le co-groupement de deux DataFrames sur la base d'une clé partagée, puis l'application d'une fonction à chaque co-groupe. La fonction reçoit deux entrées pyarrow.Table et est censée retourner une seule pyarrow.Table.

Performance

En supprimant la coûteuse conversion de données Pandas/Arrow, les UDF Arrow s'exécutent généralement plus rapidement que les UDF Pandas, avec une utilisation de mémoire réduite. Comparons les deux UDF simples :

L'UDF Arrow est environ 10 % plus rapide que l'UDF Pandas, et le profileur de mémoire montre qu'environ 40 % de mémoire est économisée lors de l'exécution.

Conclusion

Databricks Runtime 18.0 introduit les UDF Arrow natives, offrant une alternative plus rapide et plus légère aux UDF Pandas pour une exécution performante des UDF Python dans PySpark. En opérant directement sur les données Arrow et en éliminant les frais de conversion Pandas/Arrow, les UDF Arrow offrent une exécution environ 10 % plus rapide, une utilisation de mémoire environ 40 % inférieure et un meilleur support pour les types de données complexes, le tout avec une syntaxe de décorateur familière et intuitive.

Prêt à en découvrir davantage ? Essayez dès aujourd'hui les UDF Arrow natives sur Databricks dans le cadre de Databricks Runtime 18.0. Pour commencer, remplacez simplement vos UDF Pandas existantes par des UDF Arrow. Dans la plupart des cas, il suffit de quelques lignes de modification pour obtenir des gains de performance immédiats. Consultez la documentation des UDF Arrow et la documentation des UDTF pour la référence API complète et des exemples supplémentaires.

(Cet article de blog a été traduit à l'aide d'outils basés sur l'intelligence artificielle) Article original

Recevez les derniers articles dans votre boîte mail

Abonnez-vous à notre blog et recevez les derniers articles directement dans votre boîte mail.