Ir para o conteúdo principal
Código aberto

Apresentando UDFs Arrow no PySpark: Uma Substituição Mais Rápida e Eficiente para UDFs Pandas

Defina UDFs mais performáticas com facilidade.

por Ruifeng Zheng e Yicong Huang

  • Introduzimos UDFs nativas do Arrow, que operam diretamente nos dados do Arrow, eliminando a sobrecarga de conversão Pandas/Arrow em UDFs Pandas para uma execução mais rápida e menor uso de memória.
  • Também descrevemos os tipos de UDFs do Arrow para casos de uso escalar e de agregação, e UDTFs do Arrow para transformações de tabela de entrada e saída, com exemplos de código em Python e SQL.
  • Benchmarks mostram que as UDFs do Arrow são ~10% mais rápidas e usam ~40% menos memória do que as UDFs Pandas, com melhor suporte para tipos de dados complexos.

Introdução

As Funções Definidas pelo Usuário (UDFs) em Python são um mecanismo essencial de extensibilidade, mas tradicionalmente sofriam de alta sobrecarga devido à execução baseada em linha. No Apache Spark™, as UDFs Pandas resolveram parte desse problema introduzindo serialização baseada em Arrow e processamento em lote, melhorando significativamente o throughput em comparação com as UDFs Python escalares.

No entanto, as UDFs Pandas ainda apresentam limitações fundamentais:

  • A conversão de dados Pandas/Arrow introduz cópias de dados adicionais. Abordagens de cópia zero são possíveis apenas em certos casos específicos. Por exemplo, colunas com valores NULL acionarão cópias profundas.
  • Tipos de dados complexos não são bem suportados. Por exemplo, instâncias aninhadas de StructType não são suportadas para o tipo de saída em casos de uso de agregação.
Fluxos de dados da execução de UDFs Pandas no Apache Spark

Ao eliminar a conversão de dados Pandas/Arrow, as UDFs Arrow executam mais rápido que as UDFs Pandas, consomem menos memória e oferecem melhor suporte a tipos de dados.

UDFs Arrow Nativas

Temos o prazer de apresentar as UDFs Arrow Nativas a partir do Databricks Runtime 18.0 (notas de lançamento), um avanço empolgante para a execução de UDFs de alto desempenho.

As UDFs Arrow Nativas operam diretamente em dados Arrow sem converter as entradas em objetos Pandas ou NumPy. Isso preserva o layout colunar de ponta a ponta, evita cópias de dados desnecessárias e permite que as UDFs usem processamento vetorizado, aproveitando o modelo nativo de computação e memória do Arrow.

Para definir uma UDF Arrow, os usuários podem usar um novo decorador Python @arrow_udf, com tipo de retorno especificado e tipo de avaliação opcional. Por exemplo:

Os usuários também podem defini-la com o decorador existente @udf com dicas de tipo completas. Por exemplo:

Nota: A definição da função deve incluir dicas de tipo para todos os argumentos e o valor de retorno.
Este design se alinha com as interfaces das UDFs Python escalares, proporcionando uma experiência consistente e intuitiva para usuários já familiarizados com UDFs Python escalares.

O seguinte demonstra como usar a UDF Arrow:

Uso em Python:

Uso em SQL:

Oferecemos suporte para variantes de interfaces de UDFs Arrow. Incluindo Funções Escalares, Funções de Agregação e Funções de Tabela. Na API de data frame, também fornecemos mapInArrow e applyInArrow para usar UDFs Arrow. Em seguida, as apresentaremos uma a uma.

Funções Escalares Arrow

As Funções Escalares Arrow realizam transformações linha a linha. Elas são o equivalente Arrow das UDFs Pandas escalares e podem ser usadas onde uma expressão de coluna é esperada, como df.select() ou df.withColumn(). Três modos de entrada são suportados: direto, iterador e iterador de múltiplos arrays. As variantes de iterador são úteis quando a UDF requer uma inicialização única e cara (por exemplo, carregar um modelo ou compilar um padrão de regex), já que o custo de configuração é amortizado em todos os lotes. Em todos os casos, a contagem de linhas de saída deve corresponder à contagem de linhas de entrada.

  • Arrays para Array: recebendo um ou mais pyarrow.Array e retornando um pyarrow.Array. O array de entrada e saída deve ter o mesmo número de valores.
  • Iterador de Arrays para Iterador de Arrays: recebendo um iterador de pyarrow.Array e retornando um iterador de pyarrow.Array. Este tipo é útil quando a execução da UDF requer uma inicialização cara.
  • Iterador de Múltiplos Arrays para Iterador de Arrays: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um iterador de pyarrow.Array.

Funções de Agregação Arrow

As Funções de Agregação Arrow recebem uma ou mais entradas pyarrow.Array e retornam um valor escalar, reduzindo um grupo de linhas em um único resultado. Elas são o equivalente Arrow das UDFs Pandas de agregação agrupada e são usadas com groupBy().agg() ou operações de Janela. Semelhante às funções escalares, as funções de agregação também suportam três modos de entrada.

Arrays para Escalar: recebendo pyarrow.Array e retornando um valor escalar.

  • Iterador de Arrays para Escalar: recebendo um iterador de pyarrow.Array e retornando um valor escalar. Isso é útil para processar grandes volumes de dados em operações de estilo de agregação.

Iterador de Múltiplos Arrays para Escalar: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um valor escalar. Agregações mais complexas podem ser definidas.

Funções de Tabela Arrow

As Funções de Tabela Arrow, também conhecidas como UDTFs Arrow (Funções de Tabela Definidas pelo Usuário), aceitam um pyarrow.RecordBatch ou múltiplos pa.Array como entrada e produzem um pyarrow.Table como saída. Isso representa o padrão predominante para transformações de tabela-entrada, tabela-saída implementadas em Python utilizando execução colunar. As UDTFs Arrow possuem a capacidade de:

  • Retornar múltiplas colunas
  • Produzir zero, uma ou múltiplas linhas
  • Executar transformações de tabela vetorizadas empregando kernels de computação Arrow

Consequentemente, são otimamente adequadas para operações como filtragem, expansão de linhas, reestruturação de dados e geração de colunas derivadas.

A interface arrow_udtf é projetada para simplicidade, empregando uma sintaxe de decorador onde você define o tipo de retorno usando uma string formatada em DDL. Nesta configuração, o método eval recebe objetos PyArrow como entrada e espera-se que produza PyArrow Tables ou RecordBatches. A interface acomoda dois modos de entrada. Ao processar argumentos de tabela, o método eval é fornecido com um objeto pa.RecordBatch que encapsula todas as colunas da tabela de entrada:

Para argumentos escalares, o método recebe objetos pa.Array, um para cada entrada escalar:

Aqui está outro exemplo:

Esta UDTF pode funcionar de duas maneiras distintas:

Uso em Python:

Uso em SQL:

Suporte a mapInArrow e applyInArrow do DataFrame

Além das Funções Definidas pelo Usuário (UDFs) e Funções de Tabela Definidas pelo Usuário (UDTFs), o PySpark oferece APIs de Função Arrow que facilitam a aplicação direta de funções nativas Python aos dados Arrow no nível do DataFrame. Essas APIs operam de forma análoga às suas contrapartes Pandas (mapInPandas, applyInPandas), mas utilizam pyarrow.RecordBatch e pyarrow.Table em vez de Pandas DataFrames, contornando assim a sobrecarga de conversão entre os formatos Pandas e Arrow.

  • Mapear. DataFrame.mapInArrow transforma um iterador de pyarrow.RecordBatch em outro iterador de pyarrow.RecordBatch, permitindo operações no nível da linha, como filtragem, transformação ou expansão.
  • Mapeamento Agrupado. groupBy().applyInArrow() aplica uma função especificada a cada grupo, aceitando e retornando um pyarrow.Table. Esta funcionalidade é benéfica para transformações por grupo, como normalização de dados.
  • Mapeamento Co-agrupado. cogroup().applyInArrow() permite o co-agrupamento de dois DataFrames com base em uma chave compartilhada, aplicando subsequentemente uma função a cada co-grupo. A função recebe duas pyarrow.Table entradas e espera-se que retorne uma única pyarrow.Table.

Desempenho

Ao remover a cara conversão de dados Pandas/Arrow, as UDFs Arrow geralmente são executadas mais rapidamente do que as UDFs Pandas, com menor uso de memória. Vamos comparar as duas UDFs simples:

A UDF Arrow é ~10% mais rápida que a UDF Pandas, e o profiler de memória mostra que ~40% da memória é economizada na execução.

Conclusão

O Databricks Runtime 18.0 introduz as UDFs Arrow Nativas, oferecendo uma alternativa mais rápida e eficiente às UDFs Pandas para a execução de UDFs Python de alto desempenho no PySpark. Ao operar diretamente nos dados Arrow e eliminar a sobrecarga de conversão Pandas/Arrow, as UDFs Arrow proporcionam uma execução ~10% mais rápida, ~40% menos uso de memória e melhor suporte para tipos de dados complexos — tudo com uma sintaxe de decorador familiar e intuitiva.

Pronto para explorar mais? Experimente as UDFs Arrow Nativas hoje no Databricks como parte do Databricks Runtime 18.0. Para começar, basta substituir suas UDFs Pandas existentes por UDFs Arrow. Na maioria dos casos, são necessárias apenas algumas linhas de alteração para obter ganhos de desempenho imediatos. Consulte a documentação das UDFs Arrow e a documentação das UDTFs Arrow para a referência completa da API e exemplos adicionais.

(Esta publicação no blog foi traduzida utilizando ferramentas baseadas em inteligência artificial) Publicação original

Receba os posts mais recentes na sua caixa de entrada

Assine nosso blog e receba os posts mais recentes diretamente na sua caixa de entrada.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])Iterador de M\u00faltiplos Arrays para Iterador de Arrays: recebendo um iterador de uma tupla de m\u00faltiplos pyarrow.Array e retornando um iterador 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)Fun\u00e7\u00f5es de Agrega\u00e7\u00e3o ArrowAs Fun\u00e7\u00f5es de Agrega\u00e7\u00e3o Arrow recebem uma ou mais entradas pyarrow.Array e retornam um valor escalar, reduzindo um grupo de linhas em um \u00fanico resultado. Elas s\u00e3o o equivalente Arrow das UDFs Pandas de agrega\u00e7\u00e3o agrupada e s\u00e3o usadas com groupBy().agg() ou opera\u00e7\u00f5es de Janela. Semelhante \u00e0s fun\u00e7\u00f5es escalares, as fun\u00e7\u00f5es de agrega\u00e7\u00e3o tamb\u00e9m suportam tr\u00eas modos de entrada. Arrays para Escalar: recebendo pyarrow.Array e retornando um valor escalar. 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)Iterador de Arrays para Escalar: recebendo um iterador de pyarrow.Array e retornando um valor escalar. Isso \u00e9 \u00fatil para processar grandes volumes de dados em opera\u00e7\u00f5es de estilo de agrega\u00e7\u00e3o.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.0Iterador de M\u00faltiplos Arrays para Escalar: recebendo um iterador de uma tupla de m\u00faltiplos pyarrow.Array e retornando um valor escalar. Agrega\u00e7\u00f5es mais complexas podem ser definidas.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.0Fun\u00e7\u00f5es de Tabela ArrowAs Fun\u00e7\u00f5es de Tabela Arrow, tamb\u00e9m conhecidas como UDTFs Arrow (Fun\u00e7\u00f5es de Tabela Definidas pelo Usu\u00e1rio), aceitam um pyarrow.RecordBatch ou m\u00faltiplos pa.Array como entrada e produzem um pyarrow.Table como sa\u00edda. Isso representa o padr\u00e3o predominante para transforma\u00e7\u00f5es de tabela-entrada, tabela-sa\u00edda implementadas em Python utilizando execu\u00e7\u00e3o colunar. As UDTFs Arrow possuem a capacidade de:Retornar m\u00faltiplas colunasProduzir zero, uma ou m\u00faltiplas linhasExecutar transforma\u00e7\u00f5es de tabela vetorizadas empregando kernels de computa\u00e7\u00e3o ArrowConsequentemente, s\u00e3o otimamente adequadas para opera\u00e7\u00f5es como filtragem, expans\u00e3o de linhas, reestrutura\u00e7\u00e3o de dados e gera\u00e7\u00e3o de colunas derivadas.A interface arrow_udtf \u00e9 projetada para simplicidade, empregando uma sintaxe de decorador onde voc\u00ea define o tipo de retorno usando uma string formatada em DDL. Nesta configura\u00e7\u00e3o, o m\u00e9todo eval recebe objetos PyArrow como entrada e espera-se que produza PyArrow Tables ou RecordBatches. A interface acomoda dois modos de entrada. Ao processar argumentos de tabela, o m\u00e9todo eval \u00e9 fornecido com um objeto pa.RecordBatch que encapsula todas as colunas da tabela de entrada: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 resultPara argumentos escalares, o m\u00e9todo recebe objetos pa.Array, um para cada entrada escalar: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 resultAqui est\u00e1 outro exemplo: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 })Esta UDTF pode funcionar de duas maneiras distintas:Uso em 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# +----------+-------+--------+Uso em 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# +----------+-------+--------+Suporte a mapInArrow e applyInArrow do DataFrameAl\u00e9m das Fun\u00e7\u00f5es Definidas pelo Usu\u00e1rio (UDFs) e Fun\u00e7\u00f5es de Tabela Definidas pelo Usu\u00e1rio (UDTFs), o PySpark oferece APIs de Fun\u00e7\u00e3o Arrow que facilitam a aplica\u00e7\u00e3o direta de fun\u00e7\u00f5es nativas Python aos dados Arrow no n\u00edvel do DataFrame. Essas APIs operam de forma an\u00e1loga \u00e0s suas contrapartes Pandas (mapInPandas, applyInPandas), mas utilizam pyarrow.RecordBatch e pyarrow.Table em vez de Pandas DataFrames, contornando assim a sobrecarga de convers\u00e3o entre os formatos Pandas e Arrow.Mapear. DataFrame.mapInArrow transforma um iterador de pyarrow.RecordBatch em outro iterador de pyarrow.RecordBatch, permitindo opera\u00e7\u00f5es no n\u00edvel da linha, como filtragem, transforma\u00e7\u00e3o ou expans\u00e3o.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# +---+---+Mapeamento Agrupado. groupBy().applyInArrow() aplica uma fun\u00e7\u00e3o especificada a cada grupo, aceitando e retornando um pyarrow.Table. Esta funcionalidade \u00e9 ben\u00e9fica para transforma\u00e7\u00f5es por grupo, como normaliza\u00e7\u00e3o de dados.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# +---+-------------------+Mapeamento Co-agrupado. cogroup().applyInArrow() permite o co-agrupamento de dois DataFrames com base em uma chave compartilhada, aplicando subsequentemente uma fun\u00e7\u00e3o a cada co-grupo. A fun\u00e7\u00e3o recebe duas pyarrow.Table entradas e espera-se que retorne uma \u00fanica 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# +----+-----+DesempenhoAo remover a cara convers\u00e3o de dados Pandas/Arrow, as UDFs Arrow geralmente s\u00e3o executadas mais rapidamente do que as UDFs Pandas, com menor uso de mem\u00f3ria. Vamos comparar as duas UDFs 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)A UDF Arrow \u00e9 ~10% mais r\u00e1pida que a UDF Pandas, e o profiler de mem\u00f3ria mostra que ~40% da mem\u00f3ria \u00e9 economizada na execu\u00e7\u00e3o.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)Conclus\u00e3oO Databricks Runtime 18.0 introduz as UDFs Arrow Nativas, oferecendo uma alternativa mais r\u00e1pida e eficiente \u00e0s UDFs Pandas para a execu\u00e7\u00e3o de UDFs Python de alto desempenho no PySpark. Ao operar diretamente nos dados Arrow e eliminar a sobrecarga de convers\u00e3o Pandas/Arrow, as UDFs Arrow proporcionam uma execu\u00e7\u00e3o ~10% mais r\u00e1pida, ~40% menos uso de mem\u00f3ria e melhor suporte para tipos de dados complexos \u2014 tudo com uma sintaxe de decorador familiar e intuitiva.Pronto para explorar mais? Experimente as UDFs Arrow Nativas hoje no Databricks como parte do Databricks Runtime 18.0. Para come\u00e7ar, basta substituir suas UDFs Pandas existentes por UDFs Arrow. Na maioria dos casos, s\u00e3o necess\u00e1rias apenas algumas linhas de altera\u00e7\u00e3o para obter ganhos de desempenho imediatos. Consulte a documenta\u00e7\u00e3o das UDFs Arrow e a documenta\u00e7\u00e3o das UDTFs Arrow para a refer\u00eancia completa da API e exemplos adicionais.(Esta publica\u00e7\u00e3o no blog foi traduzida utilizando ferramentas baseadas em intelig\u00eancia artificial) Publica\u00e7\u00e3o original", "description": "Descubra como escrever UDFs mais perform\u00e1ticas com suporte nativo a pyarrow.", "headline": "Apresentando UDFs Arrow no PySpark: Uma Substitui\u00e7\u00e3o Mais R\u00e1pida e Eficiente para UDFs Pandas", "mainEntityOfPage": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs", "inLanguage": "pt-br", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/br/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/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "item": "https://www.databricks.com/br/blog", "name": "Todos os blogs", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "item": "https://www.databricks.com/br/blog/category/engineering", "name": "Engenharia", "position": 2}]}], "dateModified": "05/22/2026T00:00:00-08:00", "articleSection": "In\u00edcio de sess\u00e3o\nC\u00f3digo aberto", "datePublished": " 05/20/2026T00:00:00-08:00", "author": [{"@type": "Person", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/br/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/br/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"}]
Ir para o conteúdo principal
Código aberto

Apresentando UDFs Arrow no PySpark: Uma Substituição Mais Rápida e Eficiente para UDFs Pandas

Defina UDFs mais performáticas com facilidade.

por Ruifeng Zheng e Yicong Huang

  • Introduzimos UDFs nativas do Arrow, que operam diretamente nos dados do Arrow, eliminando a sobrecarga de conversão Pandas/Arrow em UDFs Pandas para uma execução mais rápida e menor uso de memória.
  • Também descrevemos os tipos de UDFs do Arrow para casos de uso escalar e de agregação, e UDTFs do Arrow para transformações de tabela de entrada e saída, com exemplos de código em Python e SQL.
  • Benchmarks mostram que as UDFs do Arrow são ~10% mais rápidas e usam ~40% menos memória do que as UDFs Pandas, com melhor suporte para tipos de dados complexos.

Introdução

As Funções Definidas pelo Usuário (UDFs) em Python são um mecanismo essencial de extensibilidade, mas tradicionalmente sofriam de alta sobrecarga devido à execução baseada em linha. No Apache Spark™, as UDFs Pandas resolveram parte desse problema introduzindo serialização baseada em Arrow e processamento em lote, melhorando significativamente o throughput em comparação com as UDFs Python escalares.

No entanto, as UDFs Pandas ainda apresentam limitações fundamentais:

  • A conversão de dados Pandas/Arrow introduz cópias de dados adicionais. Abordagens de cópia zero são possíveis apenas em certos casos específicos. Por exemplo, colunas com valores NULL acionarão cópias profundas.
  • Tipos de dados complexos não são bem suportados. Por exemplo, instâncias aninhadas de StructType não são suportadas para o tipo de saída em casos de uso de agregação.
Fluxos de dados da execução de UDFs Pandas no Apache Spark

Ao eliminar a conversão de dados Pandas/Arrow, as UDFs Arrow executam mais rápido que as UDFs Pandas, consomem menos memória e oferecem melhor suporte a tipos de dados.

UDFs Arrow Nativas

Temos o prazer de apresentar as UDFs Arrow Nativas a partir do Databricks Runtime 18.0 (notas de lançamento), um avanço empolgante para a execução de UDFs de alto desempenho.

As UDFs Arrow Nativas operam diretamente em dados Arrow sem converter as entradas em objetos Pandas ou NumPy. Isso preserva o layout colunar de ponta a ponta, evita cópias de dados desnecessárias e permite que as UDFs usem processamento vetorizado, aproveitando o modelo nativo de computação e memória do Arrow.

Para definir uma UDF Arrow, os usuários podem usar um novo decorador Python @arrow_udf, com tipo de retorno especificado e tipo de avaliação opcional. Por exemplo:

Os usuários também podem defini-la com o decorador existente @udf com dicas de tipo completas. Por exemplo:

Nota: A definição da função deve incluir dicas de tipo para todos os argumentos e o valor de retorno.
Este design se alinha com as interfaces das UDFs Python escalares, proporcionando uma experiência consistente e intuitiva para usuários já familiarizados com UDFs Python escalares.

O seguinte demonstra como usar a UDF Arrow:

Uso em Python:

Uso em SQL:

Oferecemos suporte para variantes de interfaces de UDFs Arrow. Incluindo Funções Escalares, Funções de Agregação e Funções de Tabela. Na API de data frame, também fornecemos mapInArrow e applyInArrow para usar UDFs Arrow. Em seguida, as apresentaremos uma a uma.

Funções Escalares Arrow

As Funções Escalares Arrow realizam transformações linha a linha. Elas são o equivalente Arrow das UDFs Pandas escalares e podem ser usadas onde uma expressão de coluna é esperada, como df.select() ou df.withColumn(). Três modos de entrada são suportados: direto, iterador e iterador de múltiplos arrays. As variantes de iterador são úteis quando a UDF requer uma inicialização única e cara (por exemplo, carregar um modelo ou compilar um padrão de regex), já que o custo de configuração é amortizado em todos os lotes. Em todos os casos, a contagem de linhas de saída deve corresponder à contagem de linhas de entrada.

  • Arrays para Array: recebendo um ou mais pyarrow.Array e retornando um pyarrow.Array. O array de entrada e saída deve ter o mesmo número de valores.
  • Iterador de Arrays para Iterador de Arrays: recebendo um iterador de pyarrow.Array e retornando um iterador de pyarrow.Array. Este tipo é útil quando a execução da UDF requer uma inicialização cara.
  • Iterador de Múltiplos Arrays para Iterador de Arrays: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um iterador de pyarrow.Array.

Funções de Agregação Arrow

As Funções de Agregação Arrow recebem uma ou mais entradas pyarrow.Array e retornam um valor escalar, reduzindo um grupo de linhas em um único resultado. Elas são o equivalente Arrow das UDFs Pandas de agregação agrupada e são usadas com groupBy().agg() ou operações de Janela. Semelhante às funções escalares, as funções de agregação também suportam três modos de entrada.

Arrays para Escalar: recebendo pyarrow.Array e retornando um valor escalar.

  • Iterador de Arrays para Escalar: recebendo um iterador de pyarrow.Array e retornando um valor escalar. Isso é útil para processar grandes volumes de dados em operações de estilo de agregação.

Iterador de Múltiplos Arrays para Escalar: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um valor escalar. Agregações mais complexas podem ser definidas.

Funções de Tabela Arrow

As Funções de Tabela Arrow, também conhecidas como UDTFs Arrow (Funções de Tabela Definidas pelo Usuário), aceitam um pyarrow.RecordBatch ou múltiplos pa.Array como entrada e produzem um pyarrow.Table como saída. Isso representa o padrão predominante para transformações de tabela-entrada, tabela-saída implementadas em Python utilizando execução colunar. As UDTFs Arrow possuem a capacidade de:

  • Retornar múltiplas colunas
  • Produzir zero, uma ou múltiplas linhas
  • Executar transformações de tabela vetorizadas empregando kernels de computação Arrow

Consequentemente, são otimamente adequadas para operações como filtragem, expansão de linhas, reestruturação de dados e geração de colunas derivadas.

A interface arrow_udtf é projetada para simplicidade, empregando uma sintaxe de decorador onde você define o tipo de retorno usando uma string formatada em DDL. Nesta configuração, o método eval recebe objetos PyArrow como entrada e espera-se que produza PyArrow Tables ou RecordBatches. A interface acomoda dois modos de entrada. Ao processar argumentos de tabela, o método eval é fornecido com um objeto pa.RecordBatch que encapsula todas as colunas da tabela de entrada:

Para argumentos escalares, o método recebe objetos pa.Array, um para cada entrada escalar:

Aqui está outro exemplo:

Esta UDTF pode funcionar de duas maneiras distintas:

Uso em Python:

Uso em SQL:

Suporte a mapInArrow e applyInArrow do DataFrame

Além das Funções Definidas pelo Usuário (UDFs) e Funções de Tabela Definidas pelo Usuário (UDTFs), o PySpark oferece APIs de Função Arrow que facilitam a aplicação direta de funções nativas Python aos dados Arrow no nível do DataFrame. Essas APIs operam de forma análoga às suas contrapartes Pandas (mapInPandas, applyInPandas), mas utilizam pyarrow.RecordBatch e pyarrow.Table em vez de Pandas DataFrames, contornando assim a sobrecarga de conversão entre os formatos Pandas e Arrow.

  • Mapear. DataFrame.mapInArrow transforma um iterador de pyarrow.RecordBatch em outro iterador de pyarrow.RecordBatch, permitindo operações no nível da linha, como filtragem, transformação ou expansão.
  • Mapeamento Agrupado. groupBy().applyInArrow() aplica uma função especificada a cada grupo, aceitando e retornando um pyarrow.Table. Esta funcionalidade é benéfica para transformações por grupo, como normalização de dados.
  • Mapeamento Co-agrupado. cogroup().applyInArrow() permite o co-agrupamento de dois DataFrames com base em uma chave compartilhada, aplicando subsequentemente uma função a cada co-grupo. A função recebe duas pyarrow.Table entradas e espera-se que retorne uma única pyarrow.Table.

Desempenho

Ao remover a cara conversão de dados Pandas/Arrow, as UDFs Arrow geralmente são executadas mais rapidamente do que as UDFs Pandas, com menor uso de memória. Vamos comparar as duas UDFs simples:

A UDF Arrow é ~10% mais rápida que a UDF Pandas, e o profiler de memória mostra que ~40% da memória é economizada na execução.

Conclusão

O Databricks Runtime 18.0 introduz as UDFs Arrow Nativas, oferecendo uma alternativa mais rápida e eficiente às UDFs Pandas para a execução de UDFs Python de alto desempenho no PySpark. Ao operar diretamente nos dados Arrow e eliminar a sobrecarga de conversão Pandas/Arrow, as UDFs Arrow proporcionam uma execução ~10% mais rápida, ~40% menos uso de memória e melhor suporte para tipos de dados complexos — tudo com uma sintaxe de decorador familiar e intuitiva.

Pronto para explorar mais? Experimente as UDFs Arrow Nativas hoje no Databricks como parte do Databricks Runtime 18.0. Para começar, basta substituir suas UDFs Pandas existentes por UDFs Arrow. Na maioria dos casos, são necessárias apenas algumas linhas de alteração para obter ganhos de desempenho imediatos. Consulte a documentação das UDFs Arrow e a documentação das UDTFs Arrow para a referência completa da API e exemplos adicionais.

(Esta publicação no blog foi traduzida utilizando ferramentas baseadas em inteligência artificial) Publicação original

Receba os posts mais recentes na sua caixa de entrada

Assine nosso blog e receba os posts mais recentes diretamente na sua caixa de entrada.

)\n for emails in iterator:\n yield pa.array([bool(pattern.match(e)) for e in emails.to_pylist()])Iterador de M\u00faltiplos Arrays para Iterador de Arrays: recebendo um iterador de uma tupla de m\u00faltiplos pyarrow.Array e retornando um iterador 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)Fun\u00e7\u00f5es de Agrega\u00e7\u00e3o ArrowAs Fun\u00e7\u00f5es de Agrega\u00e7\u00e3o Arrow recebem uma ou mais entradas pyarrow.Array e retornam um valor escalar, reduzindo um grupo de linhas em um \u00fanico resultado. Elas s\u00e3o o equivalente Arrow das UDFs Pandas de agrega\u00e7\u00e3o agrupada e s\u00e3o usadas com groupBy().agg() ou opera\u00e7\u00f5es de Janela. Semelhante \u00e0s fun\u00e7\u00f5es escalares, as fun\u00e7\u00f5es de agrega\u00e7\u00e3o tamb\u00e9m suportam tr\u00eas modos de entrada. Arrays para Escalar: recebendo pyarrow.Array e retornando um valor escalar. 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)Iterador de Arrays para Escalar: recebendo um iterador de pyarrow.Array e retornando um valor escalar. Isso \u00e9 \u00fatil para processar grandes volumes de dados em opera\u00e7\u00f5es de estilo de agrega\u00e7\u00e3o.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.0Iterador de M\u00faltiplos Arrays para Escalar: recebendo um iterador de uma tupla de m\u00faltiplos pyarrow.Array e retornando um valor escalar. Agrega\u00e7\u00f5es mais complexas podem ser definidas.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.0Fun\u00e7\u00f5es de Tabela ArrowAs Fun\u00e7\u00f5es de Tabela Arrow, tamb\u00e9m conhecidas como UDTFs Arrow (Fun\u00e7\u00f5es de Tabela Definidas pelo Usu\u00e1rio), aceitam um pyarrow.RecordBatch ou m\u00faltiplos pa.Array como entrada e produzem um pyarrow.Table como sa\u00edda. Isso representa o padr\u00e3o predominante para transforma\u00e7\u00f5es de tabela-entrada, tabela-sa\u00edda implementadas em Python utilizando execu\u00e7\u00e3o colunar. As UDTFs Arrow possuem a capacidade de:Retornar m\u00faltiplas colunasProduzir zero, uma ou m\u00faltiplas linhasExecutar transforma\u00e7\u00f5es de tabela vetorizadas empregando kernels de computa\u00e7\u00e3o ArrowConsequentemente, s\u00e3o otimamente adequadas para opera\u00e7\u00f5es como filtragem, expans\u00e3o de linhas, reestrutura\u00e7\u00e3o de dados e gera\u00e7\u00e3o de colunas derivadas.A interface arrow_udtf \u00e9 projetada para simplicidade, empregando uma sintaxe de decorador onde voc\u00ea define o tipo de retorno usando uma string formatada em DDL. Nesta configura\u00e7\u00e3o, o m\u00e9todo eval recebe objetos PyArrow como entrada e espera-se que produza PyArrow Tables ou RecordBatches. A interface acomoda dois modos de entrada. Ao processar argumentos de tabela, o m\u00e9todo eval \u00e9 fornecido com um objeto pa.RecordBatch que encapsula todas as colunas da tabela de entrada: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 resultPara argumentos escalares, o m\u00e9todo recebe objetos pa.Array, um para cada entrada escalar: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 resultAqui est\u00e1 outro exemplo: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 })Esta UDTF pode funcionar de duas maneiras distintas:Uso em 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# +----------+-------+--------+Uso em 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# +----------+-------+--------+Suporte a mapInArrow e applyInArrow do DataFrameAl\u00e9m das Fun\u00e7\u00f5es Definidas pelo Usu\u00e1rio (UDFs) e Fun\u00e7\u00f5es de Tabela Definidas pelo Usu\u00e1rio (UDTFs), o PySpark oferece APIs de Fun\u00e7\u00e3o Arrow que facilitam a aplica\u00e7\u00e3o direta de fun\u00e7\u00f5es nativas Python aos dados Arrow no n\u00edvel do DataFrame. Essas APIs operam de forma an\u00e1loga \u00e0s suas contrapartes Pandas (mapInPandas, applyInPandas), mas utilizam pyarrow.RecordBatch e pyarrow.Table em vez de Pandas DataFrames, contornando assim a sobrecarga de convers\u00e3o entre os formatos Pandas e Arrow.Mapear. DataFrame.mapInArrow transforma um iterador de pyarrow.RecordBatch em outro iterador de pyarrow.RecordBatch, permitindo opera\u00e7\u00f5es no n\u00edvel da linha, como filtragem, transforma\u00e7\u00e3o ou expans\u00e3o.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# +---+---+Mapeamento Agrupado. groupBy().applyInArrow() aplica uma fun\u00e7\u00e3o especificada a cada grupo, aceitando e retornando um pyarrow.Table. Esta funcionalidade \u00e9 ben\u00e9fica para transforma\u00e7\u00f5es por grupo, como normaliza\u00e7\u00e3o de dados.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# +---+-------------------+Mapeamento Co-agrupado. cogroup().applyInArrow() permite o co-agrupamento de dois DataFrames com base em uma chave compartilhada, aplicando subsequentemente uma fun\u00e7\u00e3o a cada co-grupo. A fun\u00e7\u00e3o recebe duas pyarrow.Table entradas e espera-se que retorne uma \u00fanica 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# +----+-----+DesempenhoAo remover a cara convers\u00e3o de dados Pandas/Arrow, as UDFs Arrow geralmente s\u00e3o executadas mais rapidamente do que as UDFs Pandas, com menor uso de mem\u00f3ria. Vamos comparar as duas UDFs 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)A UDF Arrow \u00e9 ~10% mais r\u00e1pida que a UDF Pandas, e o profiler de mem\u00f3ria mostra que ~40% da mem\u00f3ria \u00e9 economizada na execu\u00e7\u00e3o.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)Conclus\u00e3oO Databricks Runtime 18.0 introduz as UDFs Arrow Nativas, oferecendo uma alternativa mais r\u00e1pida e eficiente \u00e0s UDFs Pandas para a execu\u00e7\u00e3o de UDFs Python de alto desempenho no PySpark. Ao operar diretamente nos dados Arrow e eliminar a sobrecarga de convers\u00e3o Pandas/Arrow, as UDFs Arrow proporcionam uma execu\u00e7\u00e3o ~10% mais r\u00e1pida, ~40% menos uso de mem\u00f3ria e melhor suporte para tipos de dados complexos \u2014 tudo com uma sintaxe de decorador familiar e intuitiva.Pronto para explorar mais? Experimente as UDFs Arrow Nativas hoje no Databricks como parte do Databricks Runtime 18.0. Para come\u00e7ar, basta substituir suas UDFs Pandas existentes por UDFs Arrow. Na maioria dos casos, s\u00e3o necess\u00e1rias apenas algumas linhas de altera\u00e7\u00e3o para obter ganhos de desempenho imediatos. Consulte a documenta\u00e7\u00e3o das UDFs Arrow e a documenta\u00e7\u00e3o das UDTFs Arrow para a refer\u00eancia completa da API e exemplos adicionais.(Esta publica\u00e7\u00e3o no blog foi traduzida utilizando ferramentas baseadas em intelig\u00eancia artificial) Publica\u00e7\u00e3o original", "dateModified": "05/22/2026T00:00:00-08:00", "inLanguage": "pt-br", "headline": "Apresentando UDFs Arrow no PySpark: Uma Substitui\u00e7\u00e3o Mais R\u00e1pida e Eficiente para UDFs Pandas", "image": [{"@type": "ImageObject", "@id": "https://www.databricks.com/br/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"}], "datePublished": " 05/20/2026T00:00:00-08:00", "mainEntityOfPage": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs", "name": "Apresentando UDFs Arrow no PySpark: Uma Substitui\u00e7\u00e3o Mais R\u00e1pida e Eficiente para UDFs Pandas", "articleSection": "C\u00f3digo aberto", "author": [{"@type": "Person", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_0_BlogPosting_author_Person", "url": "https://www.databricks.com/br/blog/author/ruifeng-zheng", "name": "Ruifeng Zheng"}, {"@type": "Person", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight-20250224200644452_1_BlogPosting_author_Person", "url": "https://www.databricks.com/br/blog/author/yicong-huang", "name": "Yicong Huang"}], "mentions": [{"@type": "BreadcrumbList", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#BlogPosting_mentions_BreadcrumbList", "itemListElement": [{"@type": "ListItem", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_0_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Todos os blogs", "item": "https://www.databricks.com/br/blog", "position": 1}, {"@type": "ListItem", "@id": "https://www.databricks.com/br/blog/introducing-arrow-udfs-pyspark-faster-leaner-replacement-pandas-udfs#Highlight20260508174503550-32645_1_BlogPosting_mentions_BreadcrumbList_itemListElement_ListItem", "name": "Engenharia", "item": "https://www.databricks.com/br/blog/category/engineering", "position": 2}]}]}, {"@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"}]
Ir para o conteúdo principal
Código aberto

Apresentando UDFs Arrow no PySpark: Uma Substituição Mais Rápida e Eficiente para UDFs Pandas

Defina UDFs mais performáticas com facilidade.

por Ruifeng Zheng e Yicong Huang

  • Introduzimos UDFs nativas do Arrow, que operam diretamente nos dados do Arrow, eliminando a sobrecarga de conversão Pandas/Arrow em UDFs Pandas para uma execução mais rápida e menor uso de memória.
  • Também descrevemos os tipos de UDFs do Arrow para casos de uso escalar e de agregação, e UDTFs do Arrow para transformações de tabela de entrada e saída, com exemplos de código em Python e SQL.
  • Benchmarks mostram que as UDFs do Arrow são ~10% mais rápidas e usam ~40% menos memória do que as UDFs Pandas, com melhor suporte para tipos de dados complexos.

Introdução

As Funções Definidas pelo Usuário (UDFs) em Python são um mecanismo essencial de extensibilidade, mas tradicionalmente sofriam de alta sobrecarga devido à execução baseada em linha. No Apache Spark™, as UDFs Pandas resolveram parte desse problema introduzindo serialização baseada em Arrow e processamento em lote, melhorando significativamente o throughput em comparação com as UDFs Python escalares.

No entanto, as UDFs Pandas ainda apresentam limitações fundamentais:

  • A conversão de dados Pandas/Arrow introduz cópias de dados adicionais. Abordagens de cópia zero são possíveis apenas em certos casos específicos. Por exemplo, colunas com valores NULL acionarão cópias profundas.
  • Tipos de dados complexos não são bem suportados. Por exemplo, instâncias aninhadas de StructType não são suportadas para o tipo de saída em casos de uso de agregação.
Fluxos de dados da execução de UDFs Pandas no Apache Spark

Ao eliminar a conversão de dados Pandas/Arrow, as UDFs Arrow executam mais rápido que as UDFs Pandas, consomem menos memória e oferecem melhor suporte a tipos de dados.

UDFs Arrow Nativas

Temos o prazer de apresentar as UDFs Arrow Nativas a partir do Databricks Runtime 18.0 (notas de lançamento), um avanço empolgante para a execução de UDFs de alto desempenho.

As UDFs Arrow Nativas operam diretamente em dados Arrow sem converter as entradas em objetos Pandas ou NumPy. Isso preserva o layout colunar de ponta a ponta, evita cópias de dados desnecessárias e permite que as UDFs usem processamento vetorizado, aproveitando o modelo nativo de computação e memória do Arrow.

Para definir uma UDF Arrow, os usuários podem usar um novo decorador Python @arrow_udf, com tipo de retorno especificado e tipo de avaliação opcional. Por exemplo:

Os usuários também podem defini-la com o decorador existente @udf com dicas de tipo completas. Por exemplo:

Nota: A definição da função deve incluir dicas de tipo para todos os argumentos e o valor de retorno.
Este design se alinha com as interfaces das UDFs Python escalares, proporcionando uma experiência consistente e intuitiva para usuários já familiarizados com UDFs Python escalares.

O seguinte demonstra como usar a UDF Arrow:

Uso em Python:

Uso em SQL:

Oferecemos suporte para variantes de interfaces de UDFs Arrow. Incluindo Funções Escalares, Funções de Agregação e Funções de Tabela. Na API de data frame, também fornecemos mapInArrow e applyInArrow para usar UDFs Arrow. Em seguida, as apresentaremos uma a uma.

Funções Escalares Arrow

As Funções Escalares Arrow realizam transformações linha a linha. Elas são o equivalente Arrow das UDFs Pandas escalares e podem ser usadas onde uma expressão de coluna é esperada, como df.select() ou df.withColumn(). Três modos de entrada são suportados: direto, iterador e iterador de múltiplos arrays. As variantes de iterador são úteis quando a UDF requer uma inicialização única e cara (por exemplo, carregar um modelo ou compilar um padrão de regex), já que o custo de configuração é amortizado em todos os lotes. Em todos os casos, a contagem de linhas de saída deve corresponder à contagem de linhas de entrada.

  • Arrays para Array: recebendo um ou mais pyarrow.Array e retornando um pyarrow.Array. O array de entrada e saída deve ter o mesmo número de valores.
  • Iterador de Arrays para Iterador de Arrays: recebendo um iterador de pyarrow.Array e retornando um iterador de pyarrow.Array. Este tipo é útil quando a execução da UDF requer uma inicialização cara.
  • Iterador de Múltiplos Arrays para Iterador de Arrays: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um iterador de pyarrow.Array.

Funções de Agregação Arrow

As Funções de Agregação Arrow recebem uma ou mais entradas pyarrow.Array e retornam um valor escalar, reduzindo um grupo de linhas em um único resultado. Elas são o equivalente Arrow das UDFs Pandas de agregação agrupada e são usadas com groupBy().agg() ou operações de Janela. Semelhante às funções escalares, as funções de agregação também suportam três modos de entrada.

Arrays para Escalar: recebendo pyarrow.Array e retornando um valor escalar.

  • Iterador de Arrays para Escalar: recebendo um iterador de pyarrow.Array e retornando um valor escalar. Isso é útil para processar grandes volumes de dados em operações de estilo de agregação.

Iterador de Múltiplos Arrays para Escalar: recebendo um iterador de uma tupla de múltiplos pyarrow.Array e retornando um valor escalar. Agregações mais complexas podem ser definidas.

Funções de Tabela Arrow

As Funções de Tabela Arrow, também conhecidas como UDTFs Arrow (Funções de Tabela Definidas pelo Usuário), aceitam um pyarrow.RecordBatch ou múltiplos pa.Array como entrada e produzem um pyarrow.Table como saída. Isso representa o padrão predominante para transformações de tabela-entrada, tabela-saída implementadas em Python utilizando execução colunar. As UDTFs Arrow possuem a capacidade de:

  • Retornar múltiplas colunas
  • Produzir zero, uma ou múltiplas linhas
  • Executar transformações de tabela vetorizadas empregando kernels de computação Arrow

Consequentemente, são otimamente adequadas para operações como filtragem, expansão de linhas, reestruturação de dados e geração de colunas derivadas.

A interface arrow_udtf é projetada para simplicidade, empregando uma sintaxe de decorador onde você define o tipo de retorno usando uma string formatada em DDL. Nesta configuração, o método eval recebe objetos PyArrow como entrada e espera-se que produza PyArrow Tables ou RecordBatches. A interface acomoda dois modos de entrada. Ao processar argumentos de tabela, o método eval é fornecido com um objeto pa.RecordBatch que encapsula todas as colunas da tabela de entrada:

Para argumentos escalares, o método recebe objetos pa.Array, um para cada entrada escalar:

Aqui está outro exemplo:

Esta UDTF pode funcionar de duas maneiras distintas:

Uso em Python:

Uso em SQL:

Suporte a mapInArrow e applyInArrow do DataFrame

Além das Funções Definidas pelo Usuário (UDFs) e Funções de Tabela Definidas pelo Usuário (UDTFs), o PySpark oferece APIs de Função Arrow que facilitam a aplicação direta de funções nativas Python aos dados Arrow no nível do DataFrame. Essas APIs operam de forma análoga às suas contrapartes Pandas (mapInPandas, applyInPandas), mas utilizam pyarrow.RecordBatch e pyarrow.Table em vez de Pandas DataFrames, contornando assim a sobrecarga de conversão entre os formatos Pandas e Arrow.

  • Mapear. DataFrame.mapInArrow transforma um iterador de pyarrow.RecordBatch em outro iterador de pyarrow.RecordBatch, permitindo operações no nível da linha, como filtragem, transformação ou expansão.
  • Mapeamento Agrupado. groupBy().applyInArrow() aplica uma função especificada a cada grupo, aceitando e retornando um pyarrow.Table. Esta funcionalidade é benéfica para transformações por grupo, como normalização de dados.
  • Mapeamento Co-agrupado. cogroup().applyInArrow() permite o co-agrupamento de dois DataFrames com base em uma chave compartilhada, aplicando subsequentemente uma função a cada co-grupo. A função recebe duas pyarrow.Table entradas e espera-se que retorne uma única pyarrow.Table.

Desempenho

Ao remover a cara conversão de dados Pandas/Arrow, as UDFs Arrow geralmente são executadas mais rapidamente do que as UDFs Pandas, com menor uso de memória. Vamos comparar as duas UDFs simples:

A UDF Arrow é ~10% mais rápida que a UDF Pandas, e o profiler de memória mostra que ~40% da memória é economizada na execução.

Conclusão

O Databricks Runtime 18.0 introduz as UDFs Arrow Nativas, oferecendo uma alternativa mais rápida e eficiente às UDFs Pandas para a execução de UDFs Python de alto desempenho no PySpark. Ao operar diretamente nos dados Arrow e eliminar a sobrecarga de conversão Pandas/Arrow, as UDFs Arrow proporcionam uma execução ~10% mais rápida, ~40% menos uso de memória e melhor suporte para tipos de dados complexos — tudo com uma sintaxe de decorador familiar e intuitiva.

Pronto para explorar mais? Experimente as UDFs Arrow Nativas hoje no Databricks como parte do Databricks Runtime 18.0. Para começar, basta substituir suas UDFs Pandas existentes por UDFs Arrow. Na maioria dos casos, são necessárias apenas algumas linhas de alteração para obter ganhos de desempenho imediatos. Consulte a documentação das UDFs Arrow e a documentação das UDTFs Arrow para a referência completa da API e exemplos adicionais.

(Esta publicação no blog foi traduzida utilizando ferramentas baseadas em inteligência artificial) Publicação original

Receba os posts mais recentes na sua caixa de entrada

Assine nosso blog e receba os posts mais recentes diretamente na sua caixa de entrada.