SparkSession

O ponto de entrada para a programação do Spark com o conjunto de dados e a API dataframe. Um SparkSession pode ser usado para criar DataFrames, registrar DataFrames como tabelas, executar SQL em tabelas, tabelas de cache e ler arquivos parquet.

Sintaxe

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Propriedades

Propriedade Descrição
builder Interface para construir a configuração da sessão.
catalog Interface pela qual o usuário pode criar, remover, alterar ou consultar bancos de dados subjacentes, tabelas, funções etc.
client Dá acesso ao cliente Spark Connect. Apenas Spark Connect.
conf Interface de configuração de runtime para Spark.
dataSource Retorna um DataSourceRegistration para registro de fonte de dados.
profile Retorna um perfil para criação de perfil de desempenho/memória.
read Retorna um DataFrameReader que pode ser usado para ler dados como um DataFrame.
readStream Retorna um DataStreamReader que pode ser usado para ler fluxos de dados como um DataFrame de streaming.
sparkContext Retorna o SparkContext subjacente. Somente modo clássico.
streams Retorna um StreamingQueryManager que permite o gerenciamento de todas as consultas de streaming ativas.
tvf Retorna um TableValuedFunction para chamar TVFs (funções com valor de tabela).
udf Retorna uma UDFRegistration para registro UDF.
udtf Retorna uma UDTFRegistration para registro UDTF.
version A versão do Spark na qual este aplicativo está em execução.

Methods

Método Descrição
createDataFrame(data, schema, samplingRatio, verifySchema) Cria um DataFrame de um RDD, uma lista, um DataFrame pandas, um ndarray numpy ou uma tabela pyarrow.
sql(sqlQuery, args, **kwargs) Retorna um DataFrame que representa o resultado da consulta fornecida.
table(tableName) Retorna a tabela especificada como um DataFrame.
range(start, end, step, numPartitions) Cria um DataFrame com uma única coluna LongType chamada id, contendo elementos em um intervalo.
newSession() Retorna um novo SparkSession com SQLConf separado, exibições temporárias registradas e UDFs, mas SparkContext compartilhado e cache de tabela. Somente modo clássico.
getActiveSession() Retorna o SparkSession ativo para o thread atual.
active() Retorna o SparkSession ativo ou padrão para o thread atual.
stop() Interrompe o SparkContext subjacente.
addArtifacts(*path, pyfile, archive, file) Adiciona artefatos à sessão do cliente.
interruptAll() Interrompe todas as operações desta sessão atualmente em execução no servidor.
interruptTag(tag) Interrompe todas as operações desta sessão com a marca fornecida.
interruptOperation(op_id) Interrompe uma operação desta sessão com a operationId fornecida.
addTag(tag) Adiciona uma marca a ser atribuída a todas as operações iniciadas por este thread nesta sessão.
removeTag(tag) Remove uma marca adicionada anteriormente para operações iniciadas por esse thread.
getTags() Obtém as marcas atualmente definidas para serem atribuídas a todas as operações iniciadas por esse thread.
clearTags() Limpa as marcas de operação do thread atual.