latestOffset

Возвращает последнее смещение, доступное при ограничении чтения.

Смещение start можно использовать для определения количества новых данных, считываемого ограничения. Для самого первого микробатча start предоставляется из возвращаемого initialOffset()значения. Для последующих микробаток он продолжается с последнего микробатча. Источник может возвращать то же смещение, что и смещение начала, если нет данных для обработки.

ReadLimit можно использовать источником для ограничения объема возвращаемых данных. Реализуйте getDefaultReadLimit() для обеспечения правильности ReadLimit , если источник может ограничить данные на основе параметров источника.

Подсистема по-прежнему может вызываться latestOffset() , ReadAllAvailable даже если источник создает другое ограничение чтения от getDefaultReadLimit(). Источник всегда должен соблюдать заданный ReadLimit обработчиком.

Добавлено в Databricks Runtime 15.2

Синтаксис

latestOffset(start: dict, limit: ReadLimit)

Параметры

Параметр Тип Описание
start Дикт Начальная смещение микробатча для продолжения чтения.
limit ReadLimit Ограничение объема данных, возвращаемого этим вызовом.

Возвраты

dict

Дикт или рекурсивный дикт, ключ и значение которого являются примитивными типами, которые включают целочисленное, строковое и логическое значение.

Примеры

from pyspark.sql.streaming.datasource import ReadAllAvailable, ReadMaxRows

def latestOffset(self, start, limit):
    # Assume the source has 10 new records between start and latest offset
    if isinstance(limit, ReadAllAvailable):
        return {"index": start["index"] + 10}
    else:  # e.g., limit is ReadMaxRows(5)
        return {"index": start["index"] + min(10, limit.maxRows)}