clusterBy (DataStreamWriter)

Klasteruje dane wyjściowe według podanych kolumn. Rekordy z podobnymi wartościami w kolumnach klastrowania są grupowane razem w tym samym pliku. Klastrowanie zwiększa wydajność zapytań, umożliwiając wykonywanie zapytań przy użyciu predykatów w kolumnach klastrowania w celu pomijania niepotrzebnych danych. W przeciwieństwie do partycjonowania klastrowanie może być używane w kolumnach o wysokiej kardynalności.

Składnia

clusterBy(*cols)

Parametry

Parameter Typ Opis
*cols str lub list Nazwy kolumn do klastra.

Zwroty

DataStreamWriter

Examples

df = spark.readStream.format("rate").load()
df.writeStream.clusterBy("value")
# <...streaming.readwriter.DataStreamWriter object ...>

Klaster strumień źródła rate by timestamp and write to Parquet( Klaster a Rate source stream by timestamp and write to 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()