SparkSession

Het toegangspunt voor het programmeren van Spark met de Gegevensset- en DataFrame-API. Een SparkSession kan worden gebruikt voor het maken van DataFrames, het registreren van DataFrames als tabellen, het uitvoeren van SQL via tabellen, cachetabellen en het lezen van Parquet-bestanden.

Syntaxis

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

Eigenschappen

Vastgoed Beschrijving
builder Interface voor het bouwen van de sessieconfiguratie.
catalog Interface waarmee de gebruiker onderliggende databases, tabellen, functies, enzovoort kan maken, verwijderen, wijzigen of er query's op kunnen uitvoeren.
client Geeft toegang tot de Spark Connect-client. Alleen Spark Connect.
conf Runtime-configuratie-interface voor Spark.
dataSource Retourneert een DataSourceRegistration voor registratie van de gegevensbron.
profile Hiermee wordt een profiel geretourneerd voor prestatie-/geheugenprofilering.
read Retourneert een DataFrameReader die kan worden gebruikt om gegevens te lezen als een DataFrame.
readStream Hiermee wordt een DataStreamReader geretourneerd die kan worden gebruikt om gegevensstromen te lezen als een streaming DataFrame.
sparkContext Retourneert de onderliggende SparkContext. Alleen klassieke modus.
streams Hiermee wordt een StreamingQueryManager geretourneerd waarmee alle actieve streamingquery's kunnen worden beheerd.
tvf Hiermee wordt een TableValuedFunction geretourneerd voor het aanroepen van tabelwaardefuncties (TVF's).
udf Retourneert een UDFRegistration voor UDF-registratie.
udtf Retourneert een UDTFRegistration voor UDTF-registratie.
version De versie van Spark waarop deze toepassing wordt uitgevoerd.

Methods

Methode Beschrijving
createDataFrame(data, schema, samplingRatio, verifySchema) Hiermee maakt u een DataFrame op basis van een RDD, een lijst, een pandas DataFrame, een numpy ndarray of een duizendbladtabel.
sql(sqlQuery, args, **kwargs) Retourneert een DataFrame dat het resultaat van de opgegeven query vertegenwoordigt.
table(tableName) Retourneert de opgegeven tabel als een DataFrame.
range(start, end, step, numPartitions) Hiermee maakt u een DataFrame met één LongType-kolom met de naam id, die elementen in een bereik bevat.
newSession() Retourneert een nieuwe SparkSession met afzonderlijke SQLConf, geregistreerde tijdelijke weergaven en UDF's, maar gedeelde SparkContext- en tabelcache. Alleen klassieke modus.
getActiveSession() Retourneert de actieve SparkSession voor de huidige thread.
active() Retourneert de actieve of standaard SparkSession voor de huidige thread.
stop() Hiermee stopt u de onderliggende SparkContext.
addArtifacts(*path, pyfile, archive, file) Voegt artefacten toe aan de clientsessie.
interruptAll() Onderbreekt alle bewerkingen van deze sessie die momenteel op de server worden uitgevoerd.
interruptTag(tag) Onderbreekt alle bewerkingen van deze sessie met de opgegeven tag.
interruptOperation(op_id) Onderbreekt een bewerking van deze sessie met de opgegeven operationId.
addTag(tag) Hiermee voegt u een tag toe die moet worden toegewezen aan alle bewerkingen die door deze thread in deze sessie zijn gestart.
removeTag(tag) Hiermee verwijdert u een tag die eerder is toegevoegd voor bewerkingen die door deze thread zijn gestart.
getTags() Hiermee haalt u de tags op die momenteel zijn toegewezen aan alle bewerkingen die door deze thread zijn gestart.
clearTags() Hiermee wist u de bewerkingstags van de huidige thread.