Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Доступно в таблицах Delta Lake в Databricks Runtime 15.4 LTS и более поздних версиях, расширение типов позволяет изменять типы данных столбцов на более широкий тип без перезаписи файлов данных.
Все управляемые таблицы каталога Unity используют Delta Lake по умолчанию. См. управляемые таблицы Unity Catalog для Delta Lake и Apache Iceberg.
Note
Включение расширения типов приводит к обновлению протоколов чтения и записи. Это может повлиять на совместимость с внешними клиентами Delta Lake. См. сведения о совместимости функций Delta Lake и протоколах.
Таблицы с включённой функцией расширения типов можно читать только с помощью Databricks Runtime 15.4 LTS и более поздних версий.
Поддерживаемые изменения типов
Типы можно расширить в соответствии со следующими правилами:
| Тип источника | Поддержка более широкого спектра типов |
|---|---|
BYTE |
SHORT, , INTBIGINT, DECIMALDOUBLE |
SHORT |
INT, BIGINT, DECIMAL, DOUBLE |
INT |
BIGINT, DECIMAL, DOUBLE |
BIGINT |
DECIMAL |
FLOAT |
DOUBLE |
DECIMAL |
DECIMAL с большей точностью и масштабированием |
DATE |
TIMESTAMP_NTZ |
VOID |
Любой тип |
Изменения типов поддерживаются для столбцов и полей верхнего уровня, вложенных в структуры, карты и массивы.
Note
VOID для любого типа не требуется включение расширения типов в таблице. Любая операция, которая обновляет тип столбца VOID , завершается успешно без дополнительной настройки.
VOID Расширение типа доступно в Databricks Runtime 18.2 и выше.
Обработка десятичных чисел
По умолчанию Spark усечает дробную часть значения, когда операция повышает целочисленный тип до decimal или double, и при последующем процессе записи значение возвращается в целочисленный столбец. Дополнительные сведения о поведении политики назначения см. в разделе "Назначение магазина".
При изменении любого числового типа на decimalобщая точность должна быть равна или больше начальной точности. При увеличении масштаба общая точность должна увеличиваться на соответствующую сумму.
Минимальное целевое значение для типов byte, shortи int — decimal(10,0). Минимальная цель для long — decimal(20,0).
Если вы хотите добавить два десятичных разряда к полю с decimal(10,1), минимальное значение — decimal(12,3).
Включение расширения типов
Note
Включение расширения типов приводит к обновлению протоколов чтения и записи. Это может повлиять на совместимость с внешними клиентами Delta Lake. См. сведения о совместимости функций Delta Lake и протоколах.
Вы можете включить расширение типов в существующей таблице, задав для свойства таблицы delta.enableTypeWidening значение true:
ALTER TABLE <table_name> SET TBLPROPERTIES ('delta.enableTypeWidening' = 'true')
Вы также можете включить расширение типов во время создания таблицы:
CREATE TABLE T(c1 INT) TBLPROPERTIES('delta.enableTypeWidening' = 'true')
Применение изменения типа вручную
ALTER COLUMN Используйте команду для ручного изменения типов:
ALTER TABLE <table_name> ALTER COLUMN <col_name> TYPE <new_type>
Эта операция обновляет схему таблицы без перезаписи базовых файлов данных. Дополнительные сведения см. в статье ALTER TABLE.
Расширенные типы с автоматической эволюцией схемы
Используйте эволюцию схемы с расширением типов для обновления типов данных в целевых таблицах, чтобы соответствовать типу входящих данных.
Note
Без включения расширения типов, при эволюции схемы всегда предпринимаются попытки привести данные к более узким типам столбцов в целевой таблице. Если вы не хотите автоматически расширить типы данных в целевых таблицах, необходимо отключить расширение типов перед запуском рабочих нагрузок с включенной эволюцией схемы.
Чтобы использовать эволюцию схемы для расширения типа данных для столбца при его загрузке, необходимо выполнить следующие условия:
- Команда записи выполняется с включенной автоматической эволюцией схемы.
- Целевая таблица имеет поддержку расширения типа.
- Тип исходного столбца шире, чем тип целевого столбца.
- Процесс расширения типов поддерживает изменение типа данных.
Несоответствия типов, которые не соответствуют всем этим условиям, соответствуют обычным правилам применения схемы. См. соблюдение схемы.
Example
В следующих примерах показано, как расширение типов работает с эволюцией схемы.
Python
Создайте целевую таблицу со столбцом INT и исходной таблицей со столбцом BIGINT :
spark.sql("CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true')")
spark.sql("CREATE TABLE source_table (id BIGINT, data STRING)")
Используйте saveAsTable() с эволюцией схемы, чтобы автоматически расширить столбец INT до BIGINT при добавлении:
spark.table("source_table").write.mode("append").option("mergeSchema", "true").saveAsTable("target_table")
Использование MERGE INTO при эволюции схемы:
from delta.tables import DeltaTable
source_df = spark.table("source_table")
target_table = DeltaTable.forName(spark, "target_table")
(target_table.alias("target")
.merge(source_df.alias("source"), "target.id = source.id")
.withSchemaEvolution()
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute()
)
Scala
Создайте целевую таблицу со столбцом INT и исходной таблицей со столбцом BIGINT :
spark.sql("CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true')")
spark.sql("CREATE TABLE source_table (id BIGINT, data STRING)")
Используйте saveAsTable() с эволюцией схемы, чтобы автоматически расширить столбец INT до BIGINT при добавлении:
spark.table("source_table").write.mode("append").option("mergeSchema", "true").saveAsTable("target_table")
Использование MERGE INTO при эволюции схемы:
import io.delta.tables.DeltaTable
val sourceDf = spark.table("source_table")
val targetTable = DeltaTable.forName(spark, "target_table")
targetTable.alias("target")
.merge(sourceDf.alias("source"), "target.id = source.id")
.withSchemaEvolution()
.whenMatched().updateAll()
.whenNotMatched().insertAll()
.execute()
SQL
Создайте целевую таблицу со столбцом INT и исходной таблицей со столбцом BIGINT :
CREATE TABLE target_table (id INT, data STRING) TBLPROPERTIES ('delta.enableTypeWidening' = 'true');
CREATE TABLE source_table (id BIGINT, data STRING);
Используйте INSERT INTO с эволюцией схемы, чтобы автоматически расширить столбец INT до BIGINT при добавлении:
INSERT WITH SCHEMA EVOLUTION INTO target_table SELECT * FROM source_table;
Использование MERGE INTO при эволюции схемы:
MERGE WITH SCHEMA EVOLUTION INTO target_table
USING source_table
ON target_table.id = source_table.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;
Автозагрузчик
Это важно
Поддержка расширения типов в автозагрузчике доступна в общедоступной предварительной версии.
Автозагрузчик поддерживает расширение типов с помощью автоматической эволюции схемы. При использовании автозагрузчика для приема данных в таблицу Delta Lake с включенным расширением типов и эволюцией схемы типы столбцов автоматически расширяются для сопоставления входящих данных.
(spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "<path-to-schema-location>")
.load("<path-to-source-data>")
.writeStream
.option("mergeSchema", "true")
.option("checkpointLocation", "<path-to-checkpoint>")
.trigger(availableNow=True)
.toTable("table_name")
)
См. раздел "Автоматическое расширение типов с помощью автозагрузчика". Кроме того, в целевой таблице должен быть активирован режим расширения типов. См. Включить расширение типов.
Отключение функции расширения типов данных в таблице
Чтобы предотвратить расширение случайного типа в включенных таблицах, присвоив свойству значение false:
ALTER TABLE <table_name> SET TBLPROPERTIES ('delta.enableTypeWidening' = 'false')
Этот параметр предотвращает изменения будущих типов в таблице, но не удаляет функцию расширения типа или отменяет изменения предыдущего типа.
Если необходимо полностью удалить функции таблицы расширения типа, можно использовать команду DROP FEATURE, как показано в следующем примере:
ALTER TABLE <table-name> DROP FEATURE 'typeWidening' [TRUNCATE HISTORY]
Note
Для таблиц, в которых было включено расширение типов с помощью Databricks Runtime 15.4 LTS, необходимо вместо этого отключить функцию typeWidening-preview.
При удалении расширения типа Databricks перезаписывает все файлы данных, которые не соответствуют текущей схеме таблицы. См. Управление функцией удаления таблицы Delta Lake и понижение протокола таблицы.
Потоковая обработка данных из таблицы Delta Lake
Поддержка расширения типов в структурированной потоковой передаче доступна в Databricks Runtime 16.4 LTS и выше.
При потоковом чтении из таблицы Delta Lake с включенным расширением типов можно настроить автоматическое расширение типов для потоковых запросов, включив эволюцию схемы с помощью параметра mergeSchema для целевой таблицы. Целевая таблица должна иметь включенное расширение типизации. См. Включить расширение типов.
Python
(spark.readStream
.table("delta_source_table")
.writeStream
.option("checkpointLocation", "/path/to/checkpointLocation")
.option("mergeSchema", "true")
.toTable("output_table")
)
Scala
spark.readStream
.table("delta_source_table")
.writeStream
.option("checkpointLocation", "/path/to/checkpointLocation")
.option("mergeSchema", "true")
.toTable("output_table")
Если mergeSchema включен и включена поддержка расширения типов в целевой таблице:
- Изменения типов применяются автоматически к таблице на более низком уровне без ручного вмешательства.
- Новые столбцы добавляются автоматически в нижестойную схему таблицы.
Если mergeSchema не включено, значения обрабатываются в соответствии с spark.sql.storeAssignmentPolicy конфигурацией, которая по умолчанию понижает уровень значений, чтобы они соответствовали типу целевого столбца. Дополнительные сведения о поведении политики назначения см. в разделе «Назначение хранилища».
Обработка изменений типов в потоке
При потоковой передаче данных из таблицы Delta Lake можно указать путь для отслеживания схемы, чтобы отслеживать неаддитивные изменения схемы, в том числе изменения типов данных. Предоставление расположения для отслеживания схемы требуется в Databricks Runtime 18.0 и ниже, а начиная с Databricks Runtime 18.1 это необязательно.
Невозможно задать schemaTrackingLocation с помощью SQL. См. неподдерживаемые функции.
schemaTrackingLocation необходимо задать расположение в том же пути, что и контрольная точка потоковой передачи. Рассмотрим пример.
Python
checkpoint_path = "/path/to/checkpointLocation"
(spark.readStream
.option("schemaTrackingLocation", checkpoint_path)
.table("delta_source_table")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("output_table")
)
Scala
val checkpointPath = "/path/to/checkpointLocation"
spark.readStream
.option("schemaTrackingLocation", checkpointPath)
.table("delta_source_table")
.writeStream
.option("checkpointLocation", checkpointPath)
.toTable("output_table")
После указания местоположения для отслеживания схемы поток обновляет отслеживаемую схему при обнаружении изменения типа, а затем останавливается. В то время необходимо обработать изменение типа, например включение расширения типов в нижней таблице или обновление потокового запроса.
Чтобы возобновить обработку, задайте конфигурацию Spark spark.databricks.delta.streaming.allowSourceColumnTypeChange или параметр DataFrame средства чтения allowSourceColumnTypeChange, как в следующем примере:
Python
checkpoint_path = "/path/to/checkpointLocation"
(spark.readStream
.option("schemaTrackingLocation", checkpoint_path)
.option("allowSourceColumnTypeChange", "<delta_source_table_version>")
# alternatively to allow all future type changes for this stream:
# .option("allowSourceColumnTypeChange", "always")
.table("delta_source_table")
.writeStream
.option("checkpointLocation", checkpoint_path)
.toTable("output_table")
)
Scala
val checkpointPath = "/path/to/checkpointLocation"
spark.readStream
.option("schemaTrackingLocation", checkpointPath)
.option("allowSourceColumnTypeChange", "<delta_source_table_version>")
// alternatively to allow all future type changes for this stream:
// .option("allowSourceColumnTypeChange", "always")
.table("delta_source_table")
.writeStream
.option("checkpointLocation", checkpointPath)
.toTable("output_table")
SQL
-- To unblock for this particular stream just for this series of schema change(s):
SET spark.databricks.delta.streaming.allowSourceColumnTypeChange.ckpt_<checkpoint_id> = "<delta_source_table_version>"
-- To unblock for this particular stream:
SET spark.databricks.delta.streaming.allowSourceColumnTypeChange = "<delta_source_table_version>"
-- To unblock for all streams:
SET spark.databricks.delta.streaming.allowSourceColumnTypeChange = "always"
Когда поток останавливается, отображается сообщение об ошибке с идентификатором контрольной точки <checkpoint_id> и версией исходной таблицы Delta Lake <delta_source_table_version>.
Чтобы просмотреть полный список параметров потоковой передачи Delta Lake, см. раздел Delta Lake.
Конвейеры Lakeflow
Вы можете включить расширение типов для конвейеров Lakeflow на уровне конвейера или для отдельных таблиц. Расширение типов позволяет автоматически расширить типы столбцов во время выполнения конвейера без полного обновления потоковых таблиц. Изменения типов в материализованных представлениях всегда активируют полный перекомпьютер, и когда изменение типа применяется к исходной таблице, материализованные представления, зависящие от этой таблицы, требуют полного перекомпьютера, чтобы отразить новые типы.
Включение расширения типов для всего конвейера
Чтобы включить расширение типов для всех таблиц в конвейере, задайте конфигурацию pipelines.enableTypeWideningконвейера:
JSON
{
"configuration": {
"pipelines.enableTypeWidening": "true"
}
}
YAML
configuration:
pipelines.enableTypeWidening: 'true'
Включение расширения типов для определенных таблиц
Кроме того, можно включить расширение типов для отдельных таблиц, задав свойство delta.enableTypeWideningтаблицы:
Python
import dlt
@dlt.table(
table_properties={"delta.enableTypeWidening": "true"}
)
def my_table():
return spark.readStream.table("source_table")
SQL
CREATE OR REFRESH STREAMING TABLE my_table
TBLPROPERTIES ('delta.enableTypeWidening' = 'true')
AS SELECT * FROM source_table
Совместимость с последующими средствами чтения
Таблицы с включенным расширением типов можно читать только в Databricks Runtime 15.4 LTS и выше. Если вы хотите, чтобы таблица с расширенными типами в вашем конвейере была доступна для чтения в Databricks Runtime 14.3 и ниже, необходимо:
- Отключите расширение типа, удалив свойство
delta.enableTypeWidening/pipelines.enableTypeWideningили задав его значение false, и активируйте полное обновление таблицы. - Включите режим совместимости в таблице.
OpenSharing
Note
Поддержка расширения типов в OpenSharing доступна в Databricks Runtime 16.1 и выше.
Предоставление общего доступа к таблице Delta Lake с включённой функцией расширения типов поддерживается в Databricks-to-Databricks OpenSharing. Поставщик и получатель должны находиться в Databricks Runtime 16.1 или более поздней версии.
Чтобы прочитать канал данных об изменениях из таблицы Delta Lake с включенным расширением типов при использовании OpenSharing, необходимо задать для ответа формат delta:
spark.read
.format("deltaSharing")
.option("responseFormat", "delta")
.option("readChangeFeed", "true")
.option("startingVersion", "<start version>")
.option("endingVersion", "<end version>")
.load("<table>")
Чтение потока изменённых данных при изменении типов не поддерживается. Вместо этого необходимо разделить операцию на два отдельных чтения: одно заканчивается на версии таблицы, содержащей изменение типа, а другое начинается с версии, содержащей изменение типа.
Ограничения
Совместимость Apache Iceberg
Apache Iceberg не поддерживает все изменения типов, предусмотренные расширением диапазона типов. См. эволюцию схемы Айсберга.
Неподдерживаемые изменения типов включают следующие:
-
byte,shortintlongилиdecimaldouble - увеличение десятичного масштабирования
-
dateдоtimestampNTZ
Если включить UniForm с поддержкой совместимости с Iceberg для таблицы Delta Lake, применение одного из описанных выше изменений типа приведет к ошибке. См. Чтение таблиц Delta Lake клиентами Iceberg с помощью UniForm.
При применении одного из этих неподдерживаемых изменений типа в таблице Delta Lake есть два варианта:
Повторная генерация метаданных Iceberg: используйте следующую команду для их восстановления без функции расширения типа в таблице:
ALTER TABLE <table-name> SET TBLPROPERTIES ('delta.universalFormat.config.icebergCompatVersion' = '<version>')Это позволяет поддерживать единую совместимость после применения несовместимых изменений типов.
Удалите функцию расширения типа таблицы: см. раздел "Отключить функцию расширения типа таблицы".
Зависимые от типа функции
Некоторые функции SQL возвращают результаты, зависящие от входного типа данных. Например, функция hash возвращает разные хэш-значения для одного и того же логического значения, если тип аргумента отличается: hash(1::INT) возвращает другой результат, чем hash(1::BIGINT).
Другие зависимые от типа функции: xxhash64, bit_get, bit_reverse. typeof
Для стабильных результатов в запросах, использующих эти функции, необходимо явно привести значения к нужному типу:
Python
spark.read.table("table_name") \
.selectExpr("hash(CAST(column_name AS BIGINT))")
Scala
spark.read.table("main.johan_lasperas.dlt_type_widening_bronze2")
.selectExpr("hash(CAST(a AS BIGINT))")
SQL
-- Use explicit casting for stable hash values
SELECT hash(CAST(column_name AS BIGINT)) FROM table_name
Неподдерживаемые функции
- Невозможно задать расположение отслеживания схем с помощью SQL при потоковой передаче из таблицы Delta Lake с изменением типа.
- Нельзя предоставить доступ к таблице, для которой включено расширение типов, потребителям, не использующим Databricks, с помощью OpenSharing.