Aller au contenu principal

latestOffset

Retourne le décalage le plus récent disponible étant donné une limite de lecture.

Le décalage start peut être utilisé pour déterminer la quantité de nouvelles données à lire compte tenu de la limite. Pour le tout premier microbatch, start est fourni à partir de la valeur de retour de initialOffset(). Pour les microbatches suivants, il continue à partir du dernier microbatch. La source peut renvoyer le même décalage que le décalage de start s'il n'y a pas de données à traiter.

ReadLimit peut être utilisé par la source pour limiter la quantité de données renvoyées. Implémentez getDefaultReadLimit() pour fournir le ReadLimit approprié si la source peut limiter les données en fonction des options de la source.

Le moteur peut toujours appeler latestOffset() avec ReadAllAvailable même si la source produit une limite de lecture différente de getDefaultReadLimit(). La source doit toujours respecter le ReadLimit fourni par le moteur.

Ajouté dans Databricks Runtime 15.2

Syntaxe

latestOffset(start: dict, limit: ReadLimit)

parameter

parameter

Type

Description

start

dict

L'offset start du microbatch à partir duquel poursuivre la lecture.

limit

ReadLimit

La limite sur la quantité de données à renvoyer par cet appel.

parameter

Type

Description

start

dict

L'offset start du microbatch à partir duquel poursuivre la lecture.

limit

ReadLimit

La limite sur la quantité de données à renvoyer par cet appel.

Renvoie

dict

Un dictionnaire ou un dictionnaire récursif dont la clé et la valeur sont des types primitifs, ce qui inclut les entiers, les chaînes de caractères et les booléens.

Exemples

Python
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)}