將輸出依列分群。 在分群欄位中值相近的紀錄會被歸類在同一檔案中。 分群透過允許以分群欄位為謂詞的查詢跳過不必要的資料,提升查詢效率。 與分割不同,群集可用於高基數欄位。
語法
clusterBy(*cols)
參數
| 參數 | 類型 | 說明 |
|---|---|---|
*cols |
力量或列表 | 欄位名稱可以聚集。 |
退貨
DataStreamWriter
Examples
df = spark.readStream.format("rate").load()
df.writeStream.clusterBy("value")
# <...streaming.readwriter.DataStreamWriter object ...>
依時間戳將原始碼串流聚類並寫入 Parquet:
import tempfile
import time
with tempfile.TemporaryDirectory(prefix="clusterBy1") as d:
with tempfile.TemporaryDirectory(prefix="clusterBy2") as cp:
df = spark.readStream.format("rate").option("rowsPerSecond", 10).load()
q = df.writeStream.clusterBy(
"timestamp").format("parquet").option("checkpointLocation", cp).start(d)
time.sleep(5)
q.stop()
spark.read.schema(df.schema).parquet(d).show()