Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Базовый класс для записи источников данных, обрабатывающих данные с помощью PyArrow RecordBatch.
В отличие DataSourceWriterот итератора объектов Spark Row , этот класс оптимизирован для формата стрелки при записи данных. Это может повысить производительность при взаимодействии с системами или библиотеками, которые изначально поддерживают стрелку. Реализуйте этот класс и верните экземпляр, DataSource.writer() чтобы сделать источник данных доступным для записи с помощью стрелки.
Синтаксис
from pyspark.sql.datasource import DataSourceArrowWriter
class MyDataSourceArrowWriter(DataSourceArrowWriter):
def write(self, iterator):
...
Методы
| Метод | Описание |
|---|---|
write(iterator) |
Записывает итератор объектов PyArrow RecordBatch в приемник. Вызывается один раз для каждого исполнителя. Возвращает сообщение WriterCommitMessageо фиксации или None если сообщение о фиксации отсутствует. Этот метод является абстрактным и должен быть реализован. |
commit(messages) |
Фиксирует задание записи с помощью списка сообщений фиксации, собранных всеми исполнителями. Вызывается на драйвере при успешном выполнении всех задач. Наследуется от DataSourceWriter. |
abort(messages) |
Прерывает задание записи с помощью списка сообщений фиксации, собранных всеми исполнителями. Вызывается на драйвере при сбое одной или нескольких задач. Наследуется от DataSourceWriter. |
Примечания
- Драйвер собирает сообщения фиксации от всех исполнителей и передает их
commit()в случае успешного выполнения всех задач или вabort()случае сбоя любой задачи. - Если задача записи завершается ошибкой, сообщение о фиксации будет в
Noneсписке, переданномcommit()илиabort().
Примеры
Реализуйте модуль записи со стрелками, который подсчитывает строки во всех пакетах:
from dataclasses import dataclass
from pyspark.sql.datasource import DataSource, DataSourceArrowWriter, WriterCommitMessage
@dataclass
class MyCommitMessage(WriterCommitMessage):
num_rows: int
class MyDataSourceArrowWriter(DataSourceArrowWriter):
def write(self, iterator):
total_rows = 0
for batch in iterator:
total_rows += len(batch)
return MyCommitMessage(num_rows=total_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")