Класс DataFrameWriterV2

Интерфейс, используемый для записи кадра данных во внешнее хранилище с помощью API версии 2.

В большинстве случаев использования с таблицами Databricks и Delta Lake DataFrameWriterV2 предоставляет более мощные и гибкие варианты, чем исходный DataFrameWriter:

  • Улучшена поддержка свойств таблицы
  • Более точное управление секционированием
  • Возможности условной перезаписи
  • Поддержка кластеризации
  • Очистка семантики для операций создания или замены

Поддержка Spark Connect

Синтаксис

Используется DataFrame.writeTo(table) для доступа к этому интерфейсу.

Методы

Метод Описание
using(provider) Указывает поставщика для базового источника выходных данных.
option(key, value) Добавьте параметр записи. Например, чтобы создать управляемую таблицу df.writeTo("test").using("delta").option("path", "s3://test").createOrReplace():
options(**options) Добавьте параметры записи.
tableProperty(property, value) Добавьте свойство таблицы. Например, используйте tableProperty("location", "s3://test") для создания внешней (неуправляемой) таблицы.
partitionedBy(col, *cols) Секционирование выходной таблицы, созданной путем создания, создания, созданияOrReplace или замены с помощью заданных столбцов или преобразований.
clusterBy(col, *cols) Кластеризация данных по заданным столбцам для оптимизации производительности запросов.
create() Создайте таблицу из содержимого кадра данных.
replace() Замените существующую таблицу содержимым кадра данных.
createOrReplace() Создайте новую таблицу или замените существующую таблицу с содержимым кадра данных.
append() Добавьте содержимое кадра данных в выходную таблицу.
overwrite(condition) Перезаписать строки, соответствующие заданному условию фильтра, с содержимым кадра данных в выходной таблице.
overwritePartitions() Перезаписать все секции, для которых кадр данных содержит по крайней мере одну строку с содержимым кадра данных в выходной таблице.

Примеры

Создание новой таблицы

# 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()

Секционирование данных

# 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()

Настройка свойств таблицы

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

Использование параметров

# 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()

Кластеризация данных

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

Операции по замене

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

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

Операции добавления

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

Операции перезаписи

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()

Цепочка методов

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