polars.LazyFrame.sink_delta
-
Записать DataFrame в виде таблицы Delta.
Предупреждение
Эта функциональность считается нестабильной. Она может быть изменена в любой момент без того, чтобы это считалось критическим изменением.
- Параметры:
-
- target
-
URI таблицы или объект DeltaTable.
-
mode{‘error’, ‘append’, ‘overwrite’, ‘ignore’, ‘merge’} -
Способ обработки существующих данных.
- Если задано значение ‘error’, возникает ошибка, если таблица уже существует (по умолчанию).
- Если задано значение ‘append’, будут добавлены новые данные.
- Если задано значение ‘overwrite’, таблица будет заменена новыми данными.
- Если задано значение ‘ignore’, запись не будет выполняться, если таблица уже существует.
- Если задано значение ‘merge’, возвращается объект
TableMergerдля объединения данных из DataFrame с существующими данными.
- storage_options
-
Дополнительные параметры для бэкендов хранения, поддерживаемых
deltalake. Для облачных хранилищ они могут включать настройки аутентификации и т. д. - credential_provider
-
Функция, которую можно вызвать для предоставления учетных данных для облачного хранилища. Ожидается, что функция вернет словарь с ключами учетных данных и необязательным временем истечения срока их действия.
Предупреждение
Эта функциональность считается нестабильной. Она может быть изменена в любой момент без того, чтобы это считалось критическим изменением.
- delta_write_options
-
Дополнительные именованные аргументы при записи таблицы Delta Lake. Список поддерживаемых параметров записи см. здесь.
- delta_merge_options
-
Именованные аргументы, необходимые для
MERGEтаблицы Delta Lake. Список поддерживаемых параметров объединения см. здесь. - engine
-
Выберите движок для обработки запроса (по умолчанию
"auto"). Также можно передать экземплярEngine. Поддерживаются следующие имена движков:-
"auto": использовать движок, заданный с помощьюConfig.set_engine_affinityили переменной окруженияPOLARS_ENGINE_AFFINITY; если она не задана, используется"streaming". -
"in-memory": использовать перед записью движок обработки в памяти; это движок по умолчанию. -
"streaming": использовать потоковый движок, который обрабатывает запросы пакетами, снижая нагрузку на память и часто превосходя по скорости движок обработки в памяти. Вскоре он станет движком Polars по умолчанию. -
"gpu": использовать движок CUDA GPU (требуется GPU Nvidia иcudf-polars). Для точной настройки передайте объектGPUEngine.
Если выбранный движок не может выполнить запрос, Polars переключится на потоковый движок.
-
- optimizations
-
Этапы оптимизации, выполняемые при оптимизации запроса.
Предупреждение
Эта функциональность считается нестабильной. Она может быть изменена в любой момент без того, чтобы это считалось критическим изменением.
- Вызывает исключения:
-
- TypeError
-
Если DataFrame содержит неподдерживаемые типы данных.
- ArrowInvalidError
-
Если DataFrame содержит типы данных, которые невозможно привести к соответствующему примитивному типу.
- TableNotFoundError
-
Если таблица Delta не существует и выполняется действие MERGE.
Примечания
Типы данных Polars
NullиTimeне поддерживаются спецификацией протокола Delta и вызовут TypeError. При записи столбцы с типом данныхCategoricalбудут преобразованы в обычные строки (не категориальные).Столбцы Polars всегда допускают значения null. Чтобы записать данные в таблицу Delta с недопускающими null столбцами, необходимо передать пользовательскую схему pyarrow в
delta_write_options. См. последний пример ниже.Примеры
Запись набора данных, который слишком велик для размещения в памяти, в таблицу Delta Lake.
>>> lf = pl.scan_parquet( ... "/path/to/my_larger_than_ram_file.parquet" ... ) >>> table_path = "/path/to/delta-table/" >>> lf.sink_delta(table_path)
Запись DataFrame в таблицу Delta Lake в локальной файловой системе.
>>> df = pl.DataFrame( ... { ... "foo": [1, 2, 3, 4, 5], ... "bar": [6, 7, 8, 9, 10], ... "ham": ["a", "b", "c", "d", "e"], ... } ... ) >>> table_path = "/path/to/delta-table/" >>> df.lazy().sink_delta(table_path)Добавление данных в существующую таблицу Delta Lake в локальной файловой системе. Обратите внимание: операция завершится ошибкой, если схема новых данных не совпадает со схемой существующей таблицы.
>>> df.lazy().sink_delta(table_path, mode="append")
Перезапись таблицы Delta Lake с созданием новой версии. Если схемы новых и старых данных совпадают, указывать
schema_modeне требуется.>>> existing_table_path = "/path/to/delta-table/" >>> df.lazy().sink_delta( ... existing_table_path, ... mode="overwrite", ... delta_write_options={"schema_mode": "overwrite"}, ... )Запись DataFrame в таблицу Delta Lake в облачном объектном хранилище, например S3.
>>> table_path = "s3://bucket/prefix/to/delta-table/" >>> df.lazy().sink_delta( ... table_path, ... storage_options={ ... "AWS_REGION": "THE_AWS_REGION", ... "AWS_ACCESS_KEY_ID": "THE_AWS_ACCESS_KEY_ID", ... "AWS_SECRET_ACCESS_KEY": "THE_AWS_SECRET_ACCESS_KEY", ... }, ... )Запись DataFrame в таблицу Delta Lake с недопускающими null столбцами.
>>> import pyarrow as pa >>> existing_table_path = "/path/to/delta-table/" >>> df.lazy().sink_delta( ... existing_table_path, ... delta_write_options={ ... "schema": pa.schema([pa.field("foo", pa.int64(), nullable=False)]) ... }, ... )Запись DataFrame в таблицу Delta Lake со сжатием zstd. Список всех именованных аргументов
delta_write_optionsсм. в документации deltalake здесь, а сведения о свойствах Writer — в частности, здесь.>>> import deltalake >>> df.lazy().sink_delta( ... table_path, ... delta_write_options={ ... "writer_properties": deltalake.WriterProperties(compression="zstd"), ... }, ... )Объединение DataFrame с существующей таблицей Delta Lake. Список всех методов
TableMergerсм. в документации deltalake здесь.>>> df = pl.DataFrame( ... { ... "foo": [1, 2, 3, 4, 5], ... "bar": [6, 7, 8, 9, 10], ... "ham": ["a", "b", "c", "d", "e"], ... } ... ) >>> table_path = "/path/to/delta-table/" >>> ( ... df.lazy() ... .sink_delta( ... "table_path", ... mode="merge", ... delta_merge_options={ ... "predicate": "s.foo = t.foo", ... "source_alias": "s", ... "target_alias": "t", ... }, ... ) ... .when_matched_update_all() ... .when_not_matched_insert_all() ... .execute() ... )
LazyFrame.sink_delta(
target: str | Path | deltalake.DeltaTable,
*,
mode: Literal['error',
'append',
'overwrite',
'ignore',
'merge'] = 'error',
storage_options: StorageOptionsDict | None = None,
credential_provider: CredentialProviderFunction | Literal['auto'] | None = 'auto',
delta_write_options: dict[str,
Any] | None = None,
delta_merge_options: dict[str,
Any] | None = None,
engine: EngineType = 'auto',
optimizations: QueryOptFlags = (,
), ) → deltalake.table.TableMerger | None
© 2020 Ritchie Vink
© 2022 Polars contributors
Licensed under the MIT License.
https://docs.pola.rs/api/python/stable/reference/api/polars.LazyFrame.sink_delta.html