Выборочно перезаписывать данные с помощью Delta Lake

Delta Lake имеет следующие варианты отличительного выборочного перезаписывания:

Опция Сценарий использования Поддерживаемые типы вычислений Минимальная версия
REPLACE WHERE Атомарная перезапись строк, соответствующих предикату. Используйте для замен с фиксированным условием совпадения, например colA = 5 или int_col IN (1, 2, 3). Все типы вычислений. SQL в Databricks Runtime 12.2 LTS и более поздних версиях. Python и Scala в Databricks Runtime 9.1 LTS и более поздних версиях.
REPLACE USING Динамические данные перезаписываются. Заменяет все строки, соответствующие указанным столбцам, на основе сравнения значений столбцов в предоставленном наборе данных. Все типы вычислений. SQL в Databricks Runtime 16.3 и более поздних версиях. Python и Scala в Databricks Runtime 18.2 и более поздних версиях.
REPLACE ON Динамические данные перезаписываются логическим выражением. Используйте для замены сложное или безопасное условие соответствия NULL, например s.colA <=> t.colA AND s.colB <=> t.colB. Все типы вычислений. SQL в Databricks Runtime 17.1 и более поздних версиях. Python и Scala в Databricks Runtime 18.2 и более поздних версиях.
partitionOverwriteMode Устаревший режим динамической перезаписи разделов, при котором перезаписываются все существующие данные в каждом разделе, в который будут записаны новые данные. Не рекомендуется для новых рабочих нагрузок. SQL поддерживает только классические вычисления. Python и Scala поддерживают все типы вычислений. SQL, Python и Scala в Databricks Runtime 11.3 LTS и более поздних версиях.

В большинстве случаев использования Databricks рекомендует использовать REPLACE USING или REPLACE WHERE. Используйте REPLACE ON только если в вашем сценарии требуются сложные условия сопоставления или условия сопоставления с проверкой на NULL.

Дополнительные сведения о поведении замены каждого параметра см. в разделе INSERT. Полный DataFrameWriter список параметров Delta Lake см. в разделе Delta Lake и Apache Iceberg.

В Scala и Python replaceOn и replaceUsing нельзя использовать в сочетании с replaceWhere, partitionOverwriteMode или overwriteSchema.

Для пустых исходных запросов оба REPLACE USING и REPLACE ON не удаляют данные, однако, REPLACE WHERE могут удалять данные.

Important

Если данные были случайно перезаписаны, можно использовать восстановление для отмены изменения.

REPLACE WHERE

Вы можете выборочно перезаписать только данные, соответствующие произвольному выражению REPLACE WHERE.

Important

Чтобы воспользоваться добавочным обновлением при запуске REPLACE WHERE, используйте потоки REPLACE WHERE в конвейерах Lakeflow. См. раздел "Пакетная обработка с помощью потоков REPLACEWHERE".

REPLACE WHERE не требует разбиения таблицы на столбцах предиката. Предикат — это произвольное выражение над столбцами таблицы.

Чтобы атомарно заменить события января в целевой таблице на данные из replace_data:

Python

(replace_data.write
  .mode("overwrite")
  .option("replaceWhere", "start_date >= '2017-01-01' AND end_date <= '2017-01-31'")
  .saveAsTable("events")
)

Scala

replace_data.write
  .mode("overwrite")
  .option("replaceWhere", "start_date >= '2017-01-01' AND end_date <= '2017-01-31'")
  .saveAsTable("events")

SQL

INSERT INTO TABLE events REPLACE WHERE start_date >= '2017-01-01' AND end_date <= '2017-01-31' SELECT * FROM replace_data

Этот пример кода записывает данные в replace_data, проверяет соответствие всех строк предикату и выполняет атомарную замену с помощью семантики overwrite . Если какие-либо значения в операции выходят за пределы предиката, эта операция завершается ошибкой по умолчанию.

В классической вычислительной среде, чтобы изменить это поведение на overwrite значения в пределах диапазона предиката и insert записи за пределами указанного диапазона, отключите проверку ограничения, установив для spark.databricks.delta.replaceWhere.constraintCheck.enabled значение false:

Python

spark.conf.set("spark.databricks.delta.replaceWhere.constraintCheck.enabled", False)

Scala

spark.conf.set("spark.databricks.delta.replaceWhere.constraintCheck.enabled", false)

SQL

SET spark.databricks.delta.replaceWhere.constraintCheck.enabled=false

Note

REPLACE WHERE принимает boolean_expression с некоторыми ограничениями. Смотрите INSERT в справочнике по языку SQL.

Для пустых исходных запросов REPLACE WHERE может удалять строки таблицы.

Устаревшее поведение

Устаревшая версия replaceWhere доступна только для классических вычислений. Обзор классических вычислений.

При использовании устаревшего replaceWhereповедения запросы перезаписывают данные, соответствующие предикату только по столбцам секций. Следующая команда будет атомарно заменять месяц январь в целевой таблице, которая секционирована с помощью date, данными из df.

Python
(df.write
  .mode("overwrite")
  .option("replaceWhere", "birthDate >= '2017-01-01' AND birthDate <= '2017-01-31'")
  .saveAsTable("people10m")
)
Scala
df.write
  .mode("overwrite")
  .option("replaceWhere", "birthDate >= '2017-01-01' AND birthDate <= '2017-01-31'")
  .saveAsTable("people10m")

Чтобы использовать устаревшее поведение, установите spark.databricks.delta.replaceWhere.dataColumns.enabled в значение false:

Python
spark.conf.set("spark.databricks.delta.replaceWhere.dataColumns.enabled", False)
Scala
spark.conf.set("spark.databricks.delta.replaceWhere.dataColumns.enabled", false)
SQL
SET spark.databricks.delta.replaceWhere.dataColumns.enabled=false

Динамические перезаписи данных

Операции динамической перезаписи данных выборочно заменяют данные, которые соответствуют указанным ключевым столбцам или булевому выражению, оставляя все остальные данные без изменений. Поддерживаются все типы таблиц: секционированные, несекционированные и таблицы с жидкой кластеризацией.

Динамическая перезапись разделов является подмножеством поведения динамической перезаписи данных. Динамические перезаписи секций заменяют все существующие данные в каждой секции, для которых запись зафиксирует новые данные и оставляет все остальные секции без изменений. Поддерживаются только секционированные таблицы.

REPLACE USING

SQL, поддерживаемый в Databricks Runtime 16.3 и более поздней версии. Python и Scala поддерживаются в Databricks Runtime 18.2 и более поздних версиях. Различия в поведении в Databricks Runtime 16.3 и 17.1 см. в разделе "Устаревшая версия".

REPLACE USING обеспечивает независимое от вычислений, атомарное поведение перезаписи, которое работает на хранилищах SQL Databricks, бессерверных вычислениях и классических вычислениях. REPLACE USING не требуется задать конфигурацию сеанса Spark.

REPLACE USING заменяет строки, если значения в указанных столбцах равны при сравнении на равенство. Все остальные данные остаются неизменными.

Использование динамической перезаписи данных с помощью REPLACE USING:

Python

(sourceDataDF.write
  .mode("overwrite")
  .option("replaceUsing", "event_id, start_date")
  .saveAsTable("events")
)

Scala

sourceDataDF.write
  .mode("overwrite")
  .option("replaceUsing", "event_id, start_date")
  .saveAsTable("events")

SQL

INSERT INTO TABLE events
  REPLACE USING (event_id, start_date)
  SELECT * FROM source_data

Для пустых исходных REPLACE USING запросов не удаляет строки таблицы.

Для сложных или безопасных условий соответствия NULL используйте REPLACE ON вместо этого. См. REPLACE ON.

Смотрите INSERT в справочнике по языку SQL.

Устаревшее поведение

В Databricks Runtime 16.3–17.1 REPLACE USING использует устаревшее поведение и допускает только динамическую перезапись разделов, тогда как в Databricks Runtime 17.2 и более поздних версиях допускается динамическая перезапись данных.

Имейте в виду следующие ограничения и особенности REPLACE USING устаревшего поведения:

  • В предложении необходимо указать полный набор столбцов USING секционирования таблицы.
  • Всегда проверяйте, что записанные данные касаются только ожидаемых секций. Одна строка в неправильной секции может непреднамеренно перезаписать всю секцию.

REPLACE ON

SQL поддерживается в Databricks Runtime 17.1 и более поздних версиях. Python и Scala поддерживаются в Databricks Runtime 18.2 и более поздних версиях.

REPLACE ON заменяет строки, если они соответствуют заданному пользователем условию, в отличие от REPLACE USING, который заменяет строки, когда указанные столбцы считаются равными по критерию равенства. Используйте REPLACE ON , если вам нужна логика сопоставления, которая REPLACE USING не поддерживается, например, рассматривать NULL значения как равные.

При необходимости используйте targetAlias параметр, чтобы указать псевдоним для целевой таблицы и .as().alias() API, чтобы указать псевдоним для исходных данных.

Синтаксис SQL см. в разделе INSERT.

Python

(sourceDataDF.alias("s")
  .write
  .mode("overwrite")
  .option("targetAlias", "t")
  .option("replaceOn", "s.event_id <=> t.event_id AND s.start_date <=> t.start_date")
  .saveAsTable("events")
)

Scala

sourceDataDF.as("s")
  .write
  .mode("overwrite")
  .option("targetAlias", "t")
  .option("replaceOn", "s.event_id <=> t.event_id AND s.start_date <=> t.start_date")
  .saveAsTable("events")

SQL

INSERT INTO TABLE events AS t
  REPLACE ON (s.event_id <=> t.event_id AND s.start_date <=> t.start_date)
  (SELECT * FROM source_data) AS s

Для пустых исходных REPLACE ON запросов не удаляет строки таблицы.

Перезапись динамических разделов с partitionOverwriteMode (устар.)

Important

Эта функция доступна в общедоступной предварительной версии.

Databricks Runtime, начиная с версии 11.3 LTS, поддерживает динамическое переписывание секционированных таблиц в режиме перезаписи: INSERT OVERWRITE в SQL или запись в DataFrame с df.write.mode("overwrite"). Этот тип перезаписи доступен только для классических вычислений, а не для хранилищ SQL Databricks или бессерверных вычислений.

Предупреждение

По возможности используйте INSERT REPLACE USING вместо перезаписи разделов INSERT OVERWRITE PARTITION и spark.sql.sources.partitionOverwriteMode=dynamic. Перезапись раздела может использовать устаревшие данные при изменении секционирования.

Чтобы использовать режим динамической перезаписи секции, задайте для конфигурации сеанса Spark значение spark.sql.sources.partitionOverwriteModedynamic. Кроме того, можно задать опцию DataFrameWriterpartitionOverwriteMode в dynamic. Если задан вариант, специфичный для запроса, он переопределяет режим, указанный в конфигурации сеанса. Значение spark.sql.sources.partitionOverwriteMode по умолчанию — static.

В следующем примере используется partitionOverwriteMode.

SQL

SET spark.sql.sources.partitionOverwriteMode=dynamic;
INSERT OVERWRITE TABLE default.people10m SELECT * FROM morePeople;

Python

(df.write
  .mode("overwrite")
  .option("partitionOverwriteMode", "dynamic")
  .saveAsTable("default.people10m")
)

Scala

df.write
  .mode("overwrite")
  .option("partitionOverwriteMode", "dynamic")
  .saveAsTable("default.people10m")

Имейте в виду следующие ограничения и характеристики для partitionOverwriteMode:

  • Невозможно установить overwriteSchema в true.
  • Нельзя указать оба partitionOverwriteMode и replaceWhere в одной DataFrameWriter операции.
  • Если вы укажете условие replaceWhere с помощью параметра DataFrameWriter, Delta Lake применит это условие, чтобы контролировать, какие данные будут перезаписаны. Этот параметр имеет приоритет над конфигурацией уровня сеанса partitionOverwriteMode .
  • Всегда проверяйте, что записанные данные касаются только ожидаемых секций. Одна строка в неправильной секции может непреднамеренно перезаписать всю секцию.