DataSourceWriter

Базовый класс для записи источников данных.

Записи источников данных отвечают за сохранение данных в источник данных. Реализуйте этот класс и верните экземпляр, DataSource.writer() чтобы сделать источник данных доступным для записи.

Добавлено в Databricks Runtime 14.3 LTS

Синтаксис

from pyspark.sql.datasource import DataSourceWriter

class MyDataSourceWriter(DataSourceWriter):
    def write(self, iterator):
        ...

Методы

Метод Описание
write(iterator) Записывает данные в источник данных. Вызывается один раз для каждого исполнителя. Принимает итератор объектов и возвращает Rowсообщение о фиксацииWriterCommitMessage.None Этот метод является абстрактным и должен быть реализован.
commit(messages) Фиксирует задание записи с помощью списка сообщений фиксации, собранных всеми исполнителями. Вызывается на драйвере при успешном выполнении всех задач.
abort(messages) Прерывает задание записи с помощью списка сообщений фиксации, собранных всеми исполнителями. Вызывается на драйвере при сбое одной или нескольких задач.

Примечания

  • Драйвер собирает сообщения фиксации от всех исполнителей и передает их commit() в случае успешного выполнения всех задач или в abort() случае сбоя любой задачи.
  • Если задача записи завершается ошибкой, сообщение о фиксации будет в None списке, переданном commit() или abort().

Примеры

Реализуйте базовый модуль записи, который сохраняет строки в файл:

from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceWriter, WriterCommitMessage

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

class MyDataSourceWriter(DataSourceWriter):
    def __init__(self, options):
        self.path = options.get("path")

    def write(self, iterator):
        rows = list(iterator)
        with open(self.path, "w") as f:
            for row in rows:
                f.write(str(row) + "\n")
        return MyCommitMessage(num_rows=len(rows))

    def commit(self, messages):
        total = sum(m.num_rows for m in messages if m is not None)
        print(f"Committed {total} rows")

    def abort(self, messages):
        print("Write job failed, performing cleanup")