火花會議

這是使用 Dataset 和 DataFrame API 來編程 Spark 的切入點。 SparkSession 可用於建立資料框架、將資料框架註冊為資料表、在資料表上執行 SQL、快取資料表,以及讀取 parquet 檔案。

語法

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

屬性

房產 說明
builder 用於建立會話設定的介面。
catalog 使用者可透過此介面建立、刪除、修改或查詢底層資料庫、資料表、函式等。
client 提供 Spark Connect 用戶端的存取權限。 僅限 Spark Connect。
conf Spark 的執行時設定介面。
dataSource 回傳一個資料來源註冊的 DataSourceRegistration。
profile 回傳一個用於效能/記憶體分析的設定檔。
read 回傳一個 DataFrameReader,可用來以 DataFrame 讀取資料。
readStream 回傳一個 DataStreamReader,可用來以串流資料框架讀取資料串流。
sparkContext 回傳底層的 SparkContext。 僅限經典模式。
streams 回傳一個 StreamingQueryManager,允許管理所有正在進行的串流查詢。
tvf 回傳一個 TableValuedFunction 用於呼叫表值函式(TVF)。
udf 回傳 UDF 註冊的 UDFRegistration。
udtf 回傳UDTFR的UDTF註冊證。
version 這個應用程式所使用的 Spark 版本。

方法

方法 說明
createDataFrame(data, schema, samplingRatio, verifySchema) 可從 RDD、清單、pandas DataFrame、numpy ndarray 或 pyarrow 表格建立 DataFrame。
sql(sqlQuery, args, **kwargs) 回傳一個代表給定查詢結果的資料框架。
table(tableName) 回傳指定的資料表為 DataFrame。
range(start, end, step, numPartitions) 建立一個包含一個名為 id的 LongType 欄位的資料框架,該欄位包含一個範圍中的元素。
newSession() 回傳一個新的 SparkSession,擁有獨立的 SQLConf、註冊的暫存視圖和 UDF,但共享 SparkContext 和資料表快取。 僅限經典模式。
getActiveSession() 回傳目前執行緒的活躍 SparkSession。
active() 回傳目前執行緒的主動或預設 SparkSession。
stop() 停止底層的 SparkContext。
addArtifacts(*path, pyfile, archive, file) 將工件(artifact)加入客戶端會話。
interruptAll() 中斷目前伺服器上執行的該會話中的所有操作。
interruptTag(tag) 以指定標籤中斷此會話中的所有操作。
interruptOperation(op_id) 以指定操作Id中斷此會話的操作。
addTag(tag) 新增一個標籤,要指派給本執行緒在此會話中啟動的所有操作。
removeTag(tag) 移除先前為本執行緒啟動的操作新增的標籤。
getTags() 取得目前設定為指派給此執行緒啟動的所有操作的標籤。
clearTags() 清除目前執行緒的操作標籤。