Pular para o conteúdo principal

SparkSession

Ponto de partida para programação Spark com a API de conjuntos de dados e DataFrame . Uma SparkSession pode ser usada para criar DataFrames, registrar DataFrames como tabelas, executar SQL em tabelas, armazenar tabelas em cache e read.parquet arquivos .parquet.

Sintaxe

Python
from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Propriedades

Propriedade

Descrição

builder

Interface para criar a configuração da sessão.

catalog

Interface através da qual o usuário pode criar, excluir, alterar ou consultar bancos de dados, tabelas, funções, etc. subjacentes.

client

Fornece acesso ao cliente Spark Connect. Apenas Spark Connect.

conf

Interface de configuração Runtime para Spark.

dataSource

Retorna um DataSourceRegistration para registro da fonte de dados.

profile

Retorna um perfil para análise 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 dados transmitidos como um DataFrame de transmissão.

sparkContext

Retorna o SparkContext subjacente. Somente no modo clássico.

streams

Retorna um StreamingQueryManager que permite gerenciar todas as consultas de transmissão ativas.

tvf

Retorna uma TableValuedFunction para chamar funções com valor de tabela (TVFs).

udf

Retorna um UDFRegistration para registro de UDF.

udtf

Retorna um UDTFRegistration para registro de UDTF.

version

A versão do Spark na qual este aplicativo está sendo executado.

Propriedade

Descrição

builder

Interface para criar a configuração da sessão.

catalog

Interface através da qual o usuário pode criar, excluir, alterar ou consultar bancos de dados, tabelas, funções, etc. subjacentes.

client

Fornece acesso ao cliente Spark Connect. Apenas Spark Connect.

conf

Interface de configuração Runtime para Spark.

dataSource

Retorna um DataSourceRegistration para registro da fonte de dados.

profile

Retorna um perfil para análise 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 dados transmitidos como um DataFrame de transmissão.

sparkContext

Retorna o SparkContext subjacente. Somente no modo clássico.

streams

Retorna um StreamingQueryManager que permite gerenciar todas as consultas de transmissão ativas.

tvf

Retorna uma TableValuedFunction para chamar funções com valor de tabela (TVFs).

udf

Retorna um UDFRegistration para registro de UDF.

udtf

Retorna um UDTFRegistration para registro de UDTF.

version

A versão do Spark na qual este aplicativo está sendo executado.

Métodos

Método

Descrição

createDataFrame(data, schema, samplingRatio, verifySchema)

Cria um DataFrame a partir de um RDD, uma lista, um DataFrame Pandas , um ndarray do NumPy ou uma tabela do Pyarrow.

sql(sqlQuery, args, **kwargs)

Retorna um DataFrame representando 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 do tipo LongType chamada id, contendo elementos em um intervalo.

newSession()

Retorna uma nova SparkSession com SQLConf separado, visualização temporária registrada e UDFs, mas com SparkContext e cache de tabela compartilhados. Somente no modo clássico.

getActiveSession()

Retorna a SparkSession ativa para a thread atual.

active()

Retorna a SparkSession ativa ou default para a 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 que estão sendo executadas no servidor.

interruptTag(tag)

Interrompe todas as operações desta sessão com a tag especificada.

interruptOperation(op_id)

Interrompe uma operação desta sessão com o operationId fornecido.

addTag(tag)

Adiciona uma tag a ser atribuída a todas as operações iniciadas por esta thread nesta sessão.

removeTag(tag)

Remove a tag adicionada anteriormente para operações iniciadas por esta thread.

getTags()

Obtém as tags atualmente definidas para serem atribuídas a todas as operações iniciadas por esta thread.

clearTags()

Limpa tags de operações da thread atual.

Método

Descrição

createDataFrame(data, schema, samplingRatio, verifySchema)

Cria um DataFrame a partir de um RDD, uma lista, um DataFrame Pandas , um ndarray do NumPy ou uma tabela do Pyarrow.

sql(sqlQuery, args, **kwargs)

Retorna um DataFrame representando 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 do tipo LongType chamada id, contendo elementos em um intervalo.

newSession()

Retorna uma nova SparkSession com SQLConf separado, visualização temporária registrada e UDFs, mas com SparkContext e cache de tabela compartilhados. Somente no modo clássico.

getActiveSession()

Retorna a SparkSession ativa para a thread atual.

active()

Retorna a SparkSession ativa ou default para a 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 que estão sendo executadas no servidor.

interruptTag(tag)

Interrompe todas as operações desta sessão com a tag especificada.

interruptOperation(op_id)

Interrompe uma operação desta sessão com o operationId fornecido.

addTag(tag)

Adiciona uma tag a ser atribuída a todas as operações iniciadas por esta thread nesta sessão.

removeTag(tag)

Remove a tag adicionada anteriormente para operações iniciadas por esta thread.

getTags()

Obtém as tags atualmente definidas para serem atribuídas a todas as operações iniciadas por esta thread.

clearTags()

Limpa tags de operações da thread atual.