這是使用 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() |
清除目前執行緒的操作標籤。 |