write (DataSourceArrowWriter)

Записывает итератор объектов PyArrow RecordBatch в приемник.

Этот метод вызывается один раз для каждого исполнителя для записи данных в источник данных. Он принимает итератор объектов PyArrow RecordBatch и возвращает одну строку, представляющую сообщение фиксации, или None если нет сообщения фиксации.

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

Синтаксис

write(iterator: Iterator[RecordBatch])

Параметры

Параметр Тип Описание
iterator Итератор[RecordBatch] Итератор объектов PyArrow RecordBatch , представляющих входные данные.

Возвраты

WriterCommitMessage

Сериализуемое сообщение о фиксации.

Примеры

from dataclasses import dataclass

@dataclass
class MyCommitMessage(WriterCommitMessage):
    num_rows: int

def write(self, iterator: Iterator["RecordBatch"]) -> "WriterCommitMessage":
    total_rows = 0
    for batch in iterator:
        total_rows += len(batch)
    return MyCommitMessage(num_rows=total_rows)