DataFrameWriterV2 sınıfı

v2 API'sini kullanarak dış depolamaya DataFrame yazmak için kullanılan arabirim.

Databricks tabloları ve Delta Lake ile çoğu kullanım örneği için DataFrameWriterV2, özgün DataFrameWriter'dan daha güçlü ve esnek seçenekler sağlar:

  • Daha iyi tablo özelliği desteği
  • Bölümleme üzerinde daha ayrıntılı denetim
  • Koşullu üzerine yazma özellikleri
  • Kümeleme desteği
  • Oluşturma veya değiştirme işlemleri için daha net semantik

Spark Connect'i destekler

Sözdizimi

Bu arabirime erişmek için kullanın DataFrame.writeTo(table) .

Methods

Yöntem Açıklama
using(provider) Temel alınan çıktı veri kaynağı için bir sağlayıcı belirtir.
option(key, value) Yazma seçeneği ekleyin. Örneğin, yönetilen tablo oluşturmak için: df.writeTo("test").using("delta").option("path", "s3://test").createOrReplace().
options(**options) Yazma seçenekleri ekleyin.
tableProperty(property, value) Tablo özelliği ekleyin. Örneğin, bir EXTERNAL (yönetilmeyen) tablosu oluşturmak için kullanın tableProperty("location", "s3://test") .
partitionedBy(col, *cols) Verilen sütunları veya dönüşümleri kullanarak create, createOrReplace veya replace ile oluşturulan çıkış tablosunu bölümleyebilirsiniz.
clusterBy(col, *cols) Sorgu performansını iyileştirmek için verileri verilen sütunlara göre kümeler.
create() Veri çerçevesinin içeriğinden yeni bir tablo oluşturun.
replace() Var olan bir tabloyu veri çerçevesinin içeriğiyle değiştirin.
createOrReplace() Yeni bir tablo oluşturun veya mevcut bir tabloyu veri çerçevesinin içeriğiyle değiştirin.
append() Veri çerçevesinin içeriğini çıkış tablosuna ekleyin.
overwrite(condition) Verilen filtre koşuluyla çıkış tablosundaki veri çerçevesinin içeriğiyle eşleşen satırların üzerine yazın.
overwritePartitions() Veri çerçevesinin, çıktı tablosundaki veri çerçevesinin içeriğiyle en az bir satır içerdiği tüm bölümün üzerine yazın.

Örnekler

Yeni tablo oluşturma

# Create a new table with DataFrame contents
df = spark.createDataFrame([{"name": "Alice", "age": 30}])
df.writeTo("my_table").create()

# Create with a specific provider
df.writeTo("my_table").using("parquet").create()

Verileri bölümleme

# Partition by single column
df.writeTo("my_table") \
    .partitionedBy("year") \
    .create()

# Partition by multiple columns
df.writeTo("my_table") \
    .partitionedBy("year", "month") \
    .create()

# Partition using transform functions
from pyspark.sql.functions import years, months, days

df.writeTo("my_table") \
    .partitionedBy(years("date"), months("date")) \
    .create()

Tablo özelliklerini ayarlama

# Add table properties
df.writeTo("my_table") \
    .tableProperty("key1", "value1") \
    .tableProperty("key2", "value2") \
    .create()

Seçenekleri kullanma

# Add write options
df.writeTo("my_table") \
    .option("compression", "snappy") \
    .option("maxRecordsPerFile", "10000") \
    .create()

# Add multiple options at once
df.writeTo("my_table") \
    .options(compression="snappy", maxRecordsPerFile="10000") \
    .create()

Verileri kümeleme

# Cluster by columns for query optimization
df.writeTo("my_table") \
    .clusterBy("user_id", "timestamp") \
    .create()

Değiştirme işlemleri

# Replace existing table
df.writeTo("my_table") \
    .using("parquet") \
    .replace()

# Create or replace (safe operation)
df.writeTo("my_table") \
    .using("parquet") \
    .createOrReplace()

Ekleme işlemleri

# Append to existing table
df.writeTo("my_table").append()

İşlemlerin üzerine yazma

from pyspark.sql.functions import col

# Overwrite specific rows based on condition
df.writeTo("my_table") \
    .overwrite(col("date") == "2025-01-01")

# Overwrite entire partitions
df.writeTo("my_table") \
    .overwritePartitions()

Yöntem zincirleme

# Combine multiple configurations
df.writeTo("my_table") \
    .using("parquet") \
    .option("compression", "snappy") \
    .tableProperty("description", "User data table") \
    .partitionedBy("year", "month") \
    .clusterBy("user_id") \
    .createOrReplace()