Pular para o conteúdo principal

Tutorial de Apache Spark: ML com PySpark

Este tutorial de Apache Spark apresenta processamento de big data, análise e ML com PySpark.
Atualizado 17 de set. de 2026  · 15 min lido

Explorar com IA

ChatGPTClaudePerplexity

Apache Spark e Python para big data e machine learning

O Apache Spark é conhecido por ser um mecanismo rápido, fácil de usar e geral para processamento de big data, com módulos nativos para streaming, SQL, machine learning (ML) e processamento de grafos. Essa tecnologia é muito requisitada para data engineers, e data scientists também se beneficiam de aprender Spark para fazer Análise Exploratória de Dados (EDA), extração de features e, claro, ML.

Neste tutorial, você vai integrar Spark com Python por meio do PySpark, a API do Spark para Python que expõe o modelo de programação do Spark. De forma prática, você vai focar em:

Tutorial de Apache Spark

Se você prefere usar Spark com R, confira o Introduction to Spark in R with sparklyr gratuito da DataCamp ou baixe o cheat sheet de PySpark SQL.

Instalando o Apache Spark

Instalar e colocar o Spark para rodar pode ser desafiador. Nesta seção, você vai ver o passo a passo para instalar no seu computador.

O primeiro passo é conferir se você atende aos pré-requisitos. O Spark é escrito em Scala e roda em uma Java Virtual Machine (JVM). Por isso, verifique se você tem o Java Development Kit (JDK) instalado. O JDK fornece uma ou mais implementações da JVM. De preferência, use a versão mais recente; na época em que este conteúdo foi escrito, era o JDK8.

Depois disso, é hora de baixar o Spark!

Baixando pyspark com pip

Você pode baixar e instalar o PySpark com a ajuda do pip. É bem simples, como instalar qualquer outro pacote. Basta rodar o comando de sempre e deixar que o processo faça o trabalho pesado por você:

$ pip install pyspark

Como alternativa, acesse a página de downloads do Spark. Mantenha as opções padrão nas três primeiras etapas e, na etapa 4, você verá um link para download. Clique no link para baixar. Para este tutorial, baixe a versão 2.2.0 do Spark e o pacote “Pre-built for Apache Hadoop 2.7 and later”.

Observação: o download pode demorar um pouco!

Baixando o Spark com Homebrew

Você também pode instalar o Spark com o Homebrew, um gerenciador de pacotes gratuito e open source. É especialmente útil se você usa macOS.

Simplesmente rode os comandos abaixo para buscar o Spark, obter mais informações e, por fim, instalar no seu computador:

# Search for spark
$ brew search spark

# Get more information on apache-spark
$ brew info apache-spark

# Install apache-spark
$ brew install apache-spark

Baixar e configurar o Spark

Em seguida, descompacte o arquivo que apareceu na sua pasta Downloads. Isso pode acontecer automaticamente ao dar um duplo clique no arquivo spark-2.2.0-bin-hadoop2.7.tgz ou via Terminal com o comando:

$ tar xvf spark-2.2.0-bin-hadoop2.7.tgz

Depois, mova a pasta descompactada para /usr/local/spark com o comando:

$ mv spark-2.1.0-bin-hadoop2.7 /usr/local/spark

Observação: se aparecer um erro de permissão negada ao mover a pasta para o novo local, adicione sudo na frente do comando. A linha ficará $ sudo mv spark-2.1.0-bin-hadoop2.7 /usr/local/spark. Você será solicitado a inserir sua senha — geralmente a mesma que você usa para desbloquear o computador ao ligar :)

Agora que está tudo pronto, abra o arquivo README no caminho /usr/local/spark. Faça isso executando

$ cd /usr/local/spark

Isso vai levar você até a pasta certa. Em seguida, explore a pasta e leia o README incluído.

Primeiro, use $ ls para listar os arquivos e pastas dentro da pasta spark. Você verá um arquivo README.md. Abra-o com um dos comandos abaixo:

# Open and edit the file
$ nano README.md

# Just read the file 
$ cat README.md

Dica: use a tecla Tab para autocompletar enquanto digita o nome do arquivo :) Isso economiza tempo.

O README traz informações gerais sobre o Spark, documentação online, como compilar o Spark, os shells interativos de Scala e Python, programas de exemplo e muito mais.

A parte que pode mais interessar aqui é como compilar o Spark — mas isso só é relevante se você não baixou uma versão pré-compilada. Para este tutorial, você baixou uma versão pré-compilada. Pressione CTRL + X para sair do README e voltar à pasta do Spark.

Se você selecionou uma versão que ainda não foi compilada, rode o comando listado no README. Na época, era o seguinte:

$ build/mvn -DskipTests clean package run

Este comando pode levar um tempo para concluir.

Noções básicas de PySpark: RDDs

Agora que você instalou o Spark e o PySpark com sucesso, vamos começar explorando o shell interativo do Spark e consolidando alguns fundamentos necessários para dar os primeiros passos. No restante do tutorial, porém, você vai trabalhar com PySpark em um notebook Jupyter.

Aplicações Spark vs. Spark Shell

O shell interativo é um ambiente de Read-Eval(uate)-Print-Loop (REPL); ou seja, tudo o que você digita é lido, avaliado e impresso para que você continue a análise. Isso lembra o IPython, um shell interativo poderoso que você talvez conheça do Jupyter. Para saber mais, leia o post IPython or Jupyter da DataCamp.

Você pode usar o shell, disponível para Python e Scala, para todo o trabalho interativo.

Além do shell, você também pode escrever e executar aplicações Spark. Diferente das aplicações, no shell a SparkSession já está criada para você começar a trabalhar sem perder tempo com inicialização.

Talvez você esteja se perguntando: o que é a SparkSession?

É o ponto de entrada principal para a funcionalidade do Spark: representa a conexão com um cluster Spark e permite criar RDDs e transmitir variáveis nesse cluster. Ao trabalhar com Spark, tudo começa e termina na SparkSession. Observação: antes do Spark 2.0.0, os três principais objetos de conexão eram SparkContext, SqlContext e HiveContext.

Você verá mais sobre isso adiante. Por ora, foquemos no shell.

O shell do Spark em Python

Dentro da pasta spark em /usr/local/spark, rode

$ ./bin/pyspark

No início, algumas mensagens aparecerão. Depois, você verá “Spark”, assim:

Python 2.7.13 (v2.7.13:a06454b1afa1, Dec 17 2016, 12:39:47) 
[GCC 4.2.1 (Apple Inc. build 5666) (dot 3)] on darwin
Type "help", "copyright", "credits" or "license" for more information.
Using Spark's default log4j profile: org/apache/spark/log4j-defaults.properties
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
17/07/26 11:41:26 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
17/07/26 11:41:47 WARN ObjectStore: Failed to get database global_temp, returning NoSuchObjectException
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /__ / .__/\_,_/_/ /_/\_\   version 2.2.0
      /_/

Using Python version 2.7.13 (v2.7.13:a06454b1afa1, Dec 17 2016 12:39:47)
SparkSession available as 'spark'.
>>>

Quando ver isso, é sinal de que você já pode começar a experimentar no shell interativo!

Dica: se preferir usar o shell do IPython em vez do shell do Spark, defina a variável de ambiente a seguir:

export PYSPARK_DRIVER_PYTHON="/usr/local/ipython/bin/ipython"

Criando RDDs

Vamos começar pequeno e criar um RDD, o bloco de construção mais básico do Spark. Um RDD simplesmente representa dados, mas não é um único objeto, uma coleção de registros, um result set ou um dataset. Isso porque ele foi pensado para dados distribuídos em vários computadores: um único RDD pode estar espalhado por milhares de JVMs, já que o Spark particiona os dados automaticamente para obter paralelismo. Claro, você pode ajustar o paralelismo para ter mais partições. Por isso, um RDD é, na prática, uma coleção de partições.

Você cria facilmente um RDD simples com a função parallelize(), passando a ela alguns dados (um iterável, como uma lista, ou uma coleção):

>>> rdd1 = spark.sparkContext.parallelize([('a',7),('a',2),('b',2)])
>>> rdd2 = spark.sparkContext.parallelize([("a",["x","y","z"]), ("b",["p", "r"])])
>>> rdd3 = spark.sparkContext.parallelize(range(100))

Observação: o objeto SparkSession possui o SparkContext, acessível via spark.sparkContext. Por compatibilidade, ainda é possível chamar o SparkContext como sc, por exemplo rdd1 = sc.parallelize(['a',7),('a',2),('b',2)]).

Operações com RDD

Agora que você criou os RDDs, pode operar em paralelo sobre os dados distribuídos em rdd1 e rdd2. Existem dois tipos de operações: transformações e ações.

Para entender intuitivamente a diferença, algumas transformações comuns são map(), filter(), flatMap(), sample(), randomSplit(), coalesce() e repartition(); e algumas ações comuns são reduce(), collect(), first(), take(), count(), saveAsHadoopFile().

Transformações são operações “preguiçosas” sobre um RDD que geram um ou vários novos RDDs, enquanto ações produzem valores que não são RDD: retornam um result set, um número, um arquivo, …

Você pode, por exemplo, agregar todos os elementos de rdd1 usando a seguinte função lambda simples e retornar o resultado ao driver:

>>> rdd1.reduce(lambda a,b: a+b)

Executar essa linha dá o resultado: ('a', 7, 'a', 2, 'b', 2). Outro exemplo de transformação é flatMapValues(), usada em RDDs de pares chave-valor, como rdd2. Aqui você passa cada valor do RDD de pares por uma função flatMap sem alterar as chaves, como definido na lambda abaixo, e executa uma ação ao coletar os resultados com collect().

>>> rdd2.flatMapValues(lambda x: x).collect()
[('a', 'x'), ('a', 'y'), ('a', 'z'), ('b', 'p'), ('b', 'r')]

Os dados

Depois de ver o básico no shell interativo, é hora de trabalhar com dados reais. Para este tutorial, vamos usar o conjunto de dados California Housing. Vale lembrar que é um conjunto de dados “pequeno” e usar Spark aqui pode ser exagero; Este tutorial tem fins educacionais e serve para mostrar como usar PySpark para construir um modelo de machine learning.

Carregando e explorando seus dados

Mesmo sabendo um pouco mais sobre os dados, reserve um tempo para explorá-los melhor. Antes, porém, configure seu Jupyter Notebook com Spark e dê os primeiros passos para definir o SparkContext.

PySpark no Jupyter Notebook

Nesta parte do tutorial, você não usará o shell, e sim criará sua própria aplicação em um Jupyter Notebook. Você já tem tudo o que precisa instalado, então não é necessário muito para fazer o PySpark funcionar no Jupyter.

Inicie o aplicativo do notebook como sempre, rodando $ jupyter notebook. Depois, crie um novo notebook, importe a biblioteca findspark e use a função init(). Aqui, vamos passar o caminho /usr/local/spark para o init() porque temos certeza de que o Spark foi instalado ali.

# Import findspark 
import findspark

# Initialize and provide path
findspark.init("/usr/local/spark")

# Or use this alternative
#findspark.init()

Dica: se você não souber se o caminho está correto ou onde o Spark foi instalado, use findspark.find() para detectar automaticamente a localização.

Se quiser outras formas de usar Spark no Jupyter, consulte nosso Apache Spark in Python: Beginner’s Guide.

Com isso resolvido, vamos criar seu primeiro programa Spark!

Criando seu primeiro programa Spark

Primeiro, importe o SparkContext do pacote pyspark e inicialize. Lembre-se: no shell interativo do Spark você não precisou fazer isso porque ele já foi criado e inicializado automaticamente! Aqui, você fará um pouco mais de trabalho :)

Importe o módulo SparkSession de pyspark.sql e construa uma SparkSession com o método builder(). Em seguida, defina o master, o nome da aplicação, adicione algumas configurações como a memória do executor e, por fim, use getOrCreate() para obter a sessão atual do Spark ou criar uma se não houver nenhuma rodando.

# Import SparkSession
from pyspark.sql import SparkSession

# Build the SparkSession
spark = SparkSession.builder \
   .master("local") \
   .appName("Linear Regression Model") \
   .config("spark.executor.memory", "1gb") \
   .getOrCreate()
   
sc = spark.sparkContext

Observação: se você receber um FileNotFoundError do tipo: “No such file or directory: ‘/User/YourName/Downloads/spark-2.1.0-bin-hadoop2.7/./bin/spark-submit’”, será preciso (re)definir o PATH do Spark. Vá ao diretório home com $ cd e edite o arquivo .bash_profile com $ nano .bash_profile.

Adicione algo como isto ao final do arquivo

export SPARK_HOME="/usr/local/spark"

Use CTRL + X para sair e não se esqueça de salvar as alterações confirmando com Y. Depois, aplique as mudanças rodando source .bash_profile.

Dica: você também pode definir variáveis de ambiente adicionais, se quiser. Provavelmente não vai precisar, mas é útil saber. Veja alguns exemplos:

# Set a fixed value for the hash seed secret
export PYTHONHASHSEED=0

# Set an alternate Python executable
export PYSPARK_PYTHON=/usr/local/ipython/bin/ipython

# Augment the default search path for shared libraries
export LD_LIBRARY_PATH=/usr/local/ipython/bin/ipython

# Augment the default search path for private libraries 
export PYTHONPATH=$SPARK_HOME/python/lib/py4j-*-src.zip:$PYTHONPATH:$SPARK_HOME/python/

Observação: agora você inicializou uma SparkSession padrão. Na maior parte dos casos, será preciso configurá-la mais a fundo, especialmente ao trabalhar com big data. Para saber mais, veja esta página.

Carregando seus dados

Este tutorial usa o conjunto de dados California Housing. Ele apareceu no artigo de 1997 Sparse Spatial Autoregressions, de Pace, R. Kelley e Ronald Barry, publicado no periódico Statistics and Probability Letters. Os autores construíram o dataset com dados do censo da Califórnia de 1990.

Os dados contêm uma linha por grupo de bloco censitário (block group), a menor unidade geográfica para a qual o U.S. Census Bureau publica amostras (geralmente com 600 a 3.000 pessoas). Nesta amostra, um block group tem em média 1425,5 indivíduos em uma área compacta. Você encontra essas informações nesta página ou no artigo citado, disponível aqui.

Esses dados espaciais contêm 20.640 observações sobre preços de imóveis com 9 variáveis econômicas:

  • Longitude refere-se à distância angular de um local geográfico ao norte ou sul do equador para cada block group;
  • Latitude refere-se à distância angular de um local geográfico a leste ou oeste do equador para cada block group;
  • Housing median age é a idade mediana das pessoas que pertencem a um block group. Observação: a mediana é o valor que fica no ponto médio de uma distribuição de frequência;
  • Total rooms é o número total de cômodos nas casas por block group;
  • Total bedrooms é o número total de quartos nas casas por block group;
  • Population é o número de habitantes de um block group;
  • Households refere-se às unidades familiares e seus ocupantes por block group;
  • Median income registra a renda mediana das pessoas de um block group; e
  • Median house value é a variável dependente e refere-se ao valor mediano das casas por block group.

Além disso, você verá que todos os block groups com valores zero para as variáveis independentes e dependentes foram excluídos.

A median house value é a variável dependente e será o alvo (target) do seu modelo de ML.

Você pode baixar os dados aqui. Procure pela pasta houses.zip, faça o download e extraia para acessar as pastas de dados.

Depois, use o método textFile() para ler os dados da pasta onde você os baixou para RDDs. Esse método recebe um URI do arquivo — neste caso, o caminho local da sua máquina — e lê como uma coleção de linhas. Para facilitar, vamos ler não só o arquivo .data, mas também o .domain, que contém o cabeçalho. Isso permite conferir a ordem das variáveis.

# Load in the data
rdd = sc.textFile('/Users/yourName/Downloads/CaliforniaHousing/cal_housing.data')

# Load in the header
header = sc.textFile('/Users/yourName/Downloads/CaliforniaHousing/cal_housing.domain')

Exploração de dados

Você já reuniu muita informação só olhando a página do dataset, mas é sempre melhor colocar a mão na massa e inspecionar seus dados com Spark em Python.

É importante entender que, como a execução do Spark é “preguiçosa”, nada foi executado ainda. Seus dados não foram realmente lidos. As variáveis rdd e header são apenas planos de execução. Você precisa acionar o Spark; então use collect() para olhar o header:

header.collect()

O método collect() traz o RDD inteiro para uma única máquina, e você verá o seguinte resultado:

[u'longitude: continuous.', u'latitude: continuous.', u'housingMedianAge: continuous. ', u'totalRooms: continuous. ', u'totalBedrooms: continuous. ', u'population: continuous. ', u'households: continuous. ', u'medianIncome: continuous. ', u'medianHouseValue: continuous. ']

Dica: use collect() com cuidado! Isso pode estourar a memória do driver. Por isso, usar take() é mais seguro quando você quer imprimir alguns elementos do RDD. Em geral, limite o tamanho do resultado sempre que possível, como você faz com SQL.

A ordem das variáveis é a mesma que vimos acima na descrição do dataset, e todas as colunas devem ter valores contínuos. Vamos forçar o Spark a trabalhar e olhar os dados de habitação da Califórnia para confirmar.

Chame o método take() no seu RDD:

rdd.take(2)

Ao executar, você pega os 2 primeiros elementos do RDD. O resultado é o esperado: como usamos textFile(), as linhas foram lidas como strings inteiras. As entradas são separadas por vírgula, e as linhas também aparecem em sequência, separadas por vírgula:

[u'-122.230000,37.880000,41.000000,880.000000,129.000000,322.000000,126.000000,8.325200,452600.000000', u'-122.220000,37.860000,21.000000,7099.000000,1106.000000,2401.000000,1138.000000,8.301400,358500.000000']

Precisamos ajustar isso. Não é necessário dividir cada entrada em tipos ainda, mas é essencial garantir que cada linha vire um elemento separado. Para isso, use map() com uma lambda que faz split na vírgula. Depois, cheque o resultado com take() como antes:

Lembrete: funções lambda são funções anônimas criadas em tempo de execução.

# Split lines on commas
rdd = rdd.map(lambda line: line.split(","))

# Inspect the first 2 lines 
rdd.take(2)

Você verá o seguinte resultado:

[[u'-122.230000', u'37.880000', u'41.000000', u'880.000000', u'129.000000', u'322.000000', u'126.000000', u'8.325200', u'452600.000000'], [u'-122.220000', u'37.860000', u'21.000000', u'7099.000000', u'1106.000000', u'2401.000000', u'1138.000000', u'8.301400', u'358500.000000']]

Alternativamente, você pode usar estas funções para inspecionar:

# Inspect the first line 
rdd.first()

# Take top elements
rdd.top(2)

Se você está acostumado a trabalhar com Pandas ou data frames em R, talvez esperasse ver um cabeçalho — mas não há. Para facilitar a vida, vamos sair do RDD e convertê-lo para um DataFrame. Sempre que possível, prefira DataFrames a RDDs. Especialmente em Python, DataFrames têm melhor performance.

Mas qual a diferença entre os dois?

Use RDDs quando precisar de transformações e ações de baixo nível em dados não estruturados. Isso significa que você não se preocupa em impor um esquema durante o processamento nem em acessar atributos por nome/coluna. Reforçando o ponto da performance: ao usar RDDs, você abre mão dos ganhos que DataFrames oferecem para dados (semi)estruturados. Use RDDs quando quiser manipular dados com construções de programação funcional, e não com expressões de domínio específicas.

Recapitulando: vamos mudar para DataFrames para usar expressões de alto nível, fazer consultas SQL para explorar mais os dados e obter acesso colunar.

Mãos à obra.

O primeiro passo é criar um SchemaRDD, ou um RDD de objetos Row com esquema. Isso é natural, pois, assim como em um DataFrame, você quer chegar a linhas e colunas. Cada entrada fica ligada a uma linha e a uma coluna, e as colunas têm tipos.

Vamos usar map() novamente com uma lambda para mapear cada entrada a um campo de uma Row. Para visualizar, considere esta primeira linha:

[u'-122.230000', u'37.880000', u'41.000000', u'880.000000', u'129.000000', u'322.000000', u'126.000000', u'8.325200', u'452600.000000']

A lambda diz que vamos construir uma linha de um SchemaRDD e que o elemento no índice 0 terá o nome “longitude”, e assim por diante.

Com esse SchemaRDD, você converte facilmente o RDD em DataFrame com toDF().

# Import the necessary modules 
from pyspark.sql import Row

# Map the RDD to a DF
df = rdd.map(lambda line: Row(longitude=line[0], 
                              latitude=line[1], 
                              housingMedianAge=line[2],
                              totalRooms=line[3],
                              totalBedRooms=line[4],
                              population=line[5], 
                              households=line[6],
                              medianIncome=line[7],
                              medianHouseValue=line[8])).toDF()

Agora que você tem o DataFrame df, pode inspecioná-lo com os métodos já usados, como first() e take(), e também com head() e show():

# Show the top 20 rows 
df.show()

Você verá que o visual é bem diferente do RDD:

tutorial pyspark

Dica: use df.columns para retornar as colunas do DataFrame.

Os dados parecem bem organizados em colunas, mas e os tipos? Ao ler os dados, o Spark tenta inferir o esquema. Será que deu certo? Use df.dtypes ou df.printSchema() para conhecer os tipos do seu DataFrame.

# Print the data types of all `df` columns
# df.dtypes

# Print the schema of `df`
df.printSchema()

Como você não executou a primeira linha, o retorno será:

root
 |-- households: string (nullable = true)
 |-- housingMedianAge: string (nullable = true)
 |-- latitude: string (nullable = true)
 |-- longitude: string (nullable = true)
 |-- medianHouseValue: string (nullable = true)
 |-- medianIncome: string (nullable = true)
 |-- population: string (nullable = true)
 |-- totalBedRooms: string (nullable = true)
 |-- totalRooms: string (nullable = true)

Todas as colunas ainda são string… Nada bom!

Se você quiser continuar com esse DataFrame, precisa corrigir e atribuir tipos mais adequados a todas as colunas. A performance também vai agradecer. Intuitivamente, você poderia fazer algo assim, convertendo cada coluna de df para FloatType():

from pyspark.sql.types import *

df = df.withColumn("longitude", df["longitude"].cast(FloatType())) \
   .withColumn("latitude", df["latitude"].cast(FloatType())) \
   .withColumn("housingMedianAge",df["housingMedianAge"].cast(FloatType())) \
   .withColumn("totalRooms", df["totalRooms"].cast(FloatType())) \ 
   .withColumn("totalBedRooms", df["totalBedRooms"].cast(FloatType())) \ 
   .withColumn("population", df["population"].cast(FloatType())) \ 
   .withColumn("households", df["households"].cast(FloatType())) \ 
   .withColumn("medianIncome", df["medianIncome"].cast(FloatType())) \ 
   .withColumn("medianHouseValue", df["medianHouseValue"].cast(FloatType()))

Mas essas chamadas repetidas são verbosas, sujeitas a erro e não ficam bonitas. Por que não escrever uma função para fazer isso de forma mais limpa?

A função abaixo recebe um DataFrame, nomes de colunas e o novo tipo desejado. Para cada coluna, ela faz o cast para o novo tipo e retorna o DataFrame:

# Import all from `sql.types`
from pyspark.sql.types import *

# Write a custom function to convert the data type of DataFrame columns
def convertColumn(df, names, newType):
  for name in names: 
     df = df.withColumn(name, df[name].cast(newType))
  return df 

# Assign all column names to `columns`
columns = ['households', 'housingMedianAge', 'latitude', 'longitude', 'medianHouseValue', 'medianIncome', 'population', 'totalBedRooms', 'totalRooms']

# Conver the `df` columns to `FloatType()`
df = convertColumn(df, columns, FloatType())

Bem melhor! Você pode inspecionar os tipos com printSchema(), como antes.

Agora que está tudo certo, vamos de fato explorar os dados. Acesso colunar e consultas SQL são duas vantagens dos DataFrames. Começando pequeno: selecione duas colunas de df e mostre 10 linhas:

df.select('population','totalBedRooms').show(10)

Essa consulta retorna:

+----------+-------------+
|population|totalBedRooms|
+----------+-------------+
|     322.0|        129.0|
|    2401.0|       1106.0|
|     496.0|        190.0|
|     558.0|        235.0|
|     565.0|        280.0|
|     413.0|        213.0|
|    1094.0|        489.0|
|    1157.0|        687.0|
|    1206.0|        665.0|
|    1551.0|        707.0|
+----------+-------------+
only showing top 10 rows

Você também pode deixar a consulta mais complexa, como neste exemplo:

df.groupBy("housingMedianAge").count().sort("housingMedianAge",ascending=False).show()

Que retorna:

+----------------+-----+                                                        
|housingMedianAge|count|
+----------------+-----+
|            52.0| 1273|
|            51.0|   48|
|            50.0|  136|
|            49.0|  134|
|            48.0|  177|
|            47.0|  198|
|            46.0|  245|
|            45.0|  294|
|            44.0|  356|
|            43.0|  353|
|            42.0|  368|
|            41.0|  296|
|            40.0|  304|
|            39.0|  369|
|            38.0|  394|
|            37.0|  537|
|            36.0|  862|
|            35.0|  824|
|            34.0|  689|
|            33.0|  615|
+----------------+-----+
only showing top 20 rows

Além de consultar, você pode descrever os dados e obter estatísticas resumidas. Isso ajuda bastante depois!

df.describe().show()

PySpark Machine Learning

Observe os valores mínimos e máximos dos atributos numéricos. Vários têm amplitude grande: será preciso normalizar o dataset.

Pré-processamento de dados

Com o que você coletou na análise exploratória, já dá para pré-processar os dados para alimentar o modelo.

  • Não se preocupe com valores ausentes; todos os zeros foram excluídos do dataset.
  • Provavelmente, você deve padronizar os dados, já que o intervalo de mínimos e máximos é bem grande.
  • Possivelmente, há atributos adicionais que você pode criar, como quartos por cômodo ou cômodos por domicílio.
  • Sua variável dependente também é grande; para facilitar, vamos ajustar os valores levemente.

Pré-processando os valores-alvo

Começando pela medianHouseValue, sua variável dependente. Para facilitar o trabalho com o alvo, vamos expressar os valores em unidades de 100.000. Ou seja, um alvo como 452600.000000 vira 4.526:

# Import all from `sql.functions` 
from pyspark.sql.functions import *

# Adjust the values of `medianHouseValue`
df = df.withColumn("medianHouseValue", col("medianHouseValue")/100000)

# Show the first 2 lines of `df`
df.take(2)

Você vai ver claramente que os valores foram ajustados corretamente ao olhar o resultado do take():

[Row(households=126.0, housingMedianAge=41.0, latitude=37.880001068115234, longitude=-122.2300033569336, medianHouseValue=4.526, medianIncome=8.325200080871582, population=322.0, totalBedRooms=129.0, totalRooms=880.0), Row(households=1138.0, housingMedianAge=21.0, latitude=37.86000061035156, longitude=-122.22000122070312, medianHouseValue=3.585, medianIncome=8.301400184631348, population=2401.0, totalBedRooms=1106.0, totalRooms=7099.0)]

Engenharia de features

Agora que você ajustou medianHouseValue, vamos adicionar variáveis extras mencionadas acima. Vamos incluir:

  • Quartos por domicílio, número de cômodos por domicílio em cada block group;
  • População por domicílio, que indica quantas pessoas vivem em cada domicílio por block group; e
  • Quartos por cômodo, indicando a fração de quartos dentre os cômodos por block group;

Como estamos com DataFrames, use select() para selecionar as colunas com que vamos trabalhar (totalRooms, households e population). Além disso, indique que você está trabalhando com colunas usando col(). Sem isso, você não conseguirá fazer operações elemento a elemento como as divisões planejadas:

# Import all from `sql.functions` if you haven't yet
from pyspark.sql.functions import *

# Divide `totalRooms` by `households`
roomsPerHousehold = df.select(col("totalRooms")/col("households"))

# Divide `population` by `households`
populationPerHousehold = df.select(col("population")/col("households"))

# Divide `totalBedRooms` by `totalRooms`
bedroomsPerRoom = df.select(col("totalBedRooms")/col("totalRooms"))

# Add the new columns to `df`
df = df.withColumn("roomsPerHousehold", col("totalRooms")/col("households")) \
   .withColumn("populationPerHousehold", col("population")/col("households")) \
   .withColumn("bedroomsPerRoom", col("totalBedRooms")/col("totalRooms"))
   
# Inspect the result
df.first()

Você verá que, na primeira linha, há cerca de 6,98 cômodos por domicílio, os domicílios têm em média 2,5 pessoas e a fração de quartos é baixa, 0,14:

Row(households=126.0, housingMedianAge=41.0, latitude=37.880001068115234, longitude=-122.2300033569336, medianHouseValue=4.526, medianIncome=8.325200080871582, population=322.0, totalBedRooms=129.0, totalRooms=880.0, roomsPerHousehold=6.984126984126984, populationPerHousehold=2.5555555555555554, bedroomsPerRoom=0.14659090909090908)

Em seguida — já prevendo um possível problema na padronização — vamos reordenar as colunas. Como você não quer padronizar os valores da variável alvo, é melhor isolá-la.

Neste caso, use select() e passe os nomes das colunas na ordem apropriada. Vamos colocar a variável alvo medianHouseValue primeiro, para que não seja afetada pela padronização.

Observação: este também é o momento de excluir variáveis que você não quer considerar. Aqui, vamos deixar de fora longitude, latitude, housingMedianAge e totalRooms.

# Re-order and select columns
df = df.select("medianHouseValue", 
              "totalBedRooms", 
              "population", 
              "households", 
              "medianIncome", 
              "roomsPerHousehold", 
              "populationPerHousehold", 
              "bedroomsPerRoom")

Padronização

Agora que reordenamos os dados, vamos normalizá-los. Falta apenas um passo: separar as features da variável alvo. Essencialmente, precisamos isolar a primeira coluna do DataFrame do restante.

Vamos usar map() (como em RDDs) para isso. Também usaremos DenseVector(), um vetor local baseado em um array de double que representa os valores. Em outras palavras, ele armazena arrays de valores para uso no PySpark.

Depois, criamos um novo DataFrame a partir de input_data e renomeamos as colunas passando uma lista com "label" e "features":

# Import `DenseVector`
from pyspark.ml.linalg import DenseVector

# Define the `input_data` 
input_data = df.rdd.map(lambda x: (x[0], DenseVector(x[1:])))

# Replace `df` with the new DataFrame
df = spark.createDataFrame(input_data, ["label", "features"])

Agora podemos escalar os dados. Use o Spark ML: a biblioteca facilita e escala o ML em big data. Você encontra algoritmos e tudo para construir pipelines de ML. Aqui, o pré-processamento é simples e um pipeline completo seria exagero, mas, se quiser, consulte esta página.

A coluna de entrada são as features, e a coluna de saída com os valores reescalados no scaled_df se chamará "features_scaled":

# Import `StandardScaler` 
from pyspark.ml.feature import StandardScaler

# Initialize the `standardScaler`
standardScaler = StandardScaler(inputCol="features", outputCol="features_scaled")

# Fit the DataFrame to the scaler
scaler = standardScaler.fit(df)

# Transform the data in `df` with the scaler
scaled_df = scaler.transform(df)

# Inspect the result
scaled_df.take(2)

Veja o DataFrame e o resultado. Uma terceira coluna features_scaled foi adicionada, permitindo comparar com features:

[Row(label=4.526, features=DenseVector([129.0, 322.0, 126.0, 8.3252, 6.9841, 2.5556, 0.1466]), features_scaled=DenseVector([0.3062, 0.2843, 0.3296, 4.3821, 2.8228, 0.2461, 2.5264])), Row(label=3.585, features=DenseVector([1106.0, 2401.0, 1138.0, 8.3014, 6.2381, 2.1098, 0.1558]), features_scaled=DenseVector([2.6255, 2.1202, 2.9765, 4.3696, 2.5213, 0.2031, 2.6851]))]

Observação: essas linhas de código são muito parecidas com o que você faria no Scikit-Learn.

Construindo um modelo de machine learning com Spark ML

Com o pré-processamento feito, é hora de construir o modelo de Regressão Linear! Como sempre, primeiro divida os dados em treino e teste. Felizmente, isso é fácil com randomSplit():

# Split the data into train and test sets
train_data, test_data = scaled_df.randomSplit([.8,.2],seed=1234)

Você passa uma lista com dois números que representam os tamanhos desejados de treino e teste e uma semente (seed) para reprodutibilidade. Para saber mais, veja o Python Machine Learning Tutorial da DataCamp.

Sem mais delongas, vamos criar o modelo!

Observação: o argumento elasticNetParam corresponde ao α (intercepto do termo L1) e o regParam (termo de regularização) corresponde ao λ. Mais informações aqui.

# Import `LinearRegression`
from pyspark.ml.regression import LinearRegression

# Initialize `lr`
lr = LinearRegression(labelCol="label", maxIter=10, regParam=0.3, elasticNetParam=0.8)

# Fit the data to the model
linearModel = lr.fit(train_data)

Com o modelo pronto, gere previsões para os dados de teste: use transform() para prever os rótulos de test_data. Depois, use operações de RDD para extrair as previsões e os rótulos reais do DataFrame e junte os dois em uma lista chamada predictionAndLabel.

Por fim, inspecione os valores previstos e reais acessando a lista com colchetes []:

# Generate predictions
predicted = linearModel.transform(test_data)

# Extract the predictions and the "known" correct labels
predictions = predicted.select("prediction").rdd.map(lambda x: x[0])
labels = predicted.select("label").rdd.map(lambda x: x[0])

# Zip `predictions` and `labels` into a list
predictionAndLabel = predictions.zip(labels).collect()

# Print out first 5 instances of `predictionAndLabel` 
predictionAndLabel[:5]

Você verá os seguintes valores reais e previstos (nessa ordem):

[(1.4491508524918457, 0.14999), (1.5705029404692372, 0.14999), (2.148727956912464, 0.14999), (1.5831547768979277, 0.344), (1.5182107797955968, 0.398)]

Avaliando o modelo

Olhar previsões é uma coisa; melhor ainda é usar métricas para entender o quão bom está o modelo. Comece imprimindo os coeficientes e o intercepto do modelo:

# Coefficients for the model
linearModel.coefficients

# Intercept for the model
linearModel.intercept

Você terá:

# The coefficients
[0.0,0.0,0.0,0.276239709215,0.0,0.0,0.0]

# The intercept
0.990399577462

Em seguida, use o atributo summary para ver o rootMeanSquaredError e o r2:

# Get the RMSE
linearModel.summary.rootMeanSquaredError

# Get the R2
linearModel.summary.r2
  • O RMSE mede o erro entre valores previstos e observados. Quanto menor o RMSE, mais próximos estão os valores previstos dos observados.

  • O R2 (R ao quadrado), ou coeficiente de determinação, mostra quão próximos os dados estão da linha de regressão ajustada. Varia de 0 a 100% (ou 0 a 1 aqui): 0% indica que o modelo não explica a variabilidade dos dados ao redor da média; 100% indica que explica toda a variabilidade. Em geral, quanto maior o R2, melhor o ajuste do modelo.

Você verá o seguinte retorno:

# RMSE
0.8692118678997669

# R2
0.4240895287218379

Seu modelo ainda pode melhorar bastante! Se quiser continuar, brinque com os hiperparâmetros, com as variáveis incluídas no DataFrame original, etc. Mas o tutorial termina por aqui! 

Antes de ir…

Antes de encerrar, pare a SparkSession com a linha abaixo:

spark.stop()

Levando big data além

Parabéns! Você chegou ao fim deste tutorial e aprendeu a construir um modelo de regressão linear com Spark ML.

Se quiser aprender mais sobre PySpark, considere fazer o curso Introduction to PySpark da DataCamp e confira o Apache Spark Tutorial: ML with PySpark.

Tópicos
Python
Ciência de dados
Aprendizado de máquina

Aprenda mais sobre Python e PySpark

Curso

Fundamentos do PySpark

4 h
157.8K
Aprenda a implementar o gerenciamento de dados distribuídos e o machine learning no Spark usando o pacote PySpark.
Ver detalhesRight Arrow
Iniciar Curso
Ver maisRight Arrow
Relacionado

Tutorial

Tutorial do Pyspark: Primeiros passos com o Pyspark

Descubra o que é o Pyspark e como ele pode ser usado, com exemplos.
Natassha Selvaraj's photo

Natassha Selvaraj

10 min

Clustering k-means

Tutorial

Introdução ao k-Means Clustering com o scikit-learn em Python

Neste tutorial, saiba como aplicar o k-Means Clustering com o scikit-learn em Python

Kevin Babitz

8 min

Tutorial

Tutorial de manipulação de dados categóricos de aprendizado de máquina com Python

Aprenda os truques comuns para lidar com dados categóricos e pré-processá-los para criar modelos de aprendizado de máquina!
Moez Ali's photo

Moez Ali

14 min

Tutorial

Como treinar um LLM com o PyTorch

Domine o processo de treinamento de grandes modelos de linguagem usando o PyTorch, desde a configuração inicial até a implementação final.
Zoumana Keita 's photo

Zoumana Keita

8 min

Tutorial

Tutorial de Python: Streamlit

Este tutorial sobre o Streamlit foi criado para ajudar cientistas de dados ou engenheiros de aprendizado de máquina que não são desenvolvedores da Web e não estão interessados em passar semanas aprendendo a usar essas estruturas para criar aplicativos da Web.
Nadia mhadhbi's photo

Nadia mhadhbi

15 min

Tutorial

Tutorial de regressão Lasso e Ridge em Python

Saiba mais sobre as técnicas de regressão lasso e ridge. Compare e analise os métodos em detalhes.
DataCamp Team's photo

DataCamp Team

10 min

Ver MaisVer Mais