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 |
|---|---|---|
| dict | L'offset start du microbatch à partir duquel poursuivre la lecture. |
| 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
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)}