Spec-Zone.ru › Polars

polars.LazyFrame.sink_delta

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

Записать DataFrame в виде таблицы Delta.

движок:ПотоковыйРаспределённый

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

Эта функциональность считается нестабильной. Она может быть изменена в любой момент без того, чтобы это считалось критическим изменением.

Параметры:
target

URI таблицы или объект DeltaTable.

mode{‘error’, ‘append’, ‘overwrite’, ‘ignore’, ‘merge’}

Способ обработки существующих данных.

  • Если задано значение ‘error’, возникает ошибка, если таблица уже существует (по умолчанию).
  • Если задано значение ‘append’, будут добавлены новые данные.
  • Если задано значение ‘overwrite’, таблица будет заменена новыми данными.
  • Если задано значение ‘ignore’, запись не будет выполняться, если таблица уже существует.
  • Если задано значение ‘merge’, возвращается объект TableMerger для объединения данных из DataFrame с существующими данными.
storage_options

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

  • Список поддерживаемых параметров хранения для S3 см. здесь.
  • Список поддерживаемых параметров хранения для GCS см. здесь.
  • Список поддерживаемых параметров хранения для Azure см. здесь.
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()
... )  

© 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

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API