easyfabric.loaders.file_tracker
get_tracker_file_path
def get_tracker_file_path(bronze_abfs_base_folder: str) -> str
Returns the ABFS path for the tracker file inside the table's Bronze folder.
The tracker always lives at::
<bronze_abfs_base_folder>/_tracking/tracker.json
For azblob (date-partitioned) pass the folder above the date partition,
e.g. abfss://…/Files/afas/Projecten.
For fabricfiles pass the folder that directly contains the files,
e.g. abfss://…/Files/afas/Projecten.
Arguments:
bronze_abfs_base_folder- Non-date-partitioned Bronze folder for this table.
Returns:
str- Full ABFS path to_tracking/tracker.json.
as_utc
def as_utc(moment: datetime) -> datetime
Returns moment in UTC, reading a naive datetime as UTC.
Sources differ: the Azure SDK hands back timezone-aware UTC on
blob.last_modified, while tracker entries written before the loader
moved to UTC, and the date-folder fallback in the bronze cleanup, carry
naive values that were always meant as UTC.
tracker_size
def tracker_size(tracker_path: str) -> int | None
Returns the tracker's size in bytes, or None when it cannot be read.
The size is what makes a short read detectable: asking how many bytes were returned says nothing without knowing how many there were.
read_whole_tracker
def read_whole_tracker(tracker_path: str) -> str
Returns the tracker's content, read in full where the runtime allows it.
notebookutils.fs.head honours a larger limit on Python compute, but the
Spark runtime caps it at its own default whatever is passed: a 142.252-byte
tracker came back as exactly 102.400 bytes there. Spark reads the file
itself and has no such cap, so it is used whenever a session is active —
the same escape logging_utils already takes for log files.
looks_truncated
def looks_truncated(content: str, size: int | None) -> bool
Reports whether less was read than the file holds.
One byte of slack: Spark hands back the lines without the file's trailing newline, which must not read as a truncated file. An unknown size cannot prove anything, so it counts as not truncated — the parse result still has the last word.
read_tracker
def read_tracker(tracker_path: str) -> tuple[list[dict], bool]
Reads the tracker (NDJSON — one entry per line) and reports whether every byte was read and every line parsed.
Reading it whole is not a given: the Spark runtime caps
notebookutils.fs.head at 100 KB regardless of the limit passed, which
hid every entry past that point from the comparison and left a half record
at the cut — the Skipping unparseable tracker line warning that started
this. :func:read_whole_tracker goes through Spark where it can, and the
result is checked against the real file size rather than against the limit
that was requested.
Arguments:
tracker_path- ABFS path returned by :func:get_tracker_file_path.
Returns:
The entry dicts oldest-first, together with True only when the whole
file was read and parsed. A trim must never rewrite the tracker on an
incomplete read — it would discard whatever it could not see.
load_previous_snapshot
def load_previous_snapshot(tracker_path: str, shortcode: str) -> list[dict]
Loads this object's historical entries from the tracker (NDJSON — one entry per line).
One tracker file can hold more than one object: the fabricfiles path is derived from the Bronze folder and the source folder, so two objects reading the same source share it. Every entry records the object that wrote it, and reading returns only those — without the filter the second object compared against the first object's snapshot, concluded the source was unchanged and skipped its own load without an error (Bug #469).
The comparison ignores case, so re-casing an object name keeps its own history readable instead of orphaning it. The Generator's duplicate-tracker rule compares paths the same way, so two objects that differ only in case are one object to both halves.
Entries without a shortcode come from before the tracker recorded one and cannot be attributed, so they are ignored. That makes the first run after the upgrade reload once, which is recoverable; counting them would keep the silent skip alive, which is not.
Returns an empty list on first run (file does not exist yet).
Arguments:
tracker_path- ABFS path returned by :func:get_tracker_file_path.shortcode- Dataplatform object name (table_config.dataplatformobjectname), as written by :func:save_snapshot.
Returns:
List of entry dicts ordered oldest-first.
save_snapshot
def save_snapshot(tracker_path: str,
file_list: list[ObjectInfo],
shortcode: str,
existing: list[dict] | None = None,
max_entries: int = 1000) -> None
Appends the current run's file metadata to the tracker (NDJSON format) and
trims the log to max_entries total lines when the cap is exceeded.
Uses notebookutils.fs.append() for normal writes — the same approach as
OneLakeFileHandler — so no read-before-write is needed on the happy path.
put() is only used during a trim to rewrite the compacted file.
Arguments:
-
tracker_path- ABFS path from :func:get_tracker_file_path. -
file_list- Current :class:ObjectInfolist frompull_files. -
shortcode- Dataplatform object name (table_config.dataplatformobjectname). -
existing- Previously loaded snapshot (from :func:load_previous_snapshot). Used only to decide whether the cap can have been reached; the trim itself always re-reads, so it rewrites the file as it is after the append rather than as it was before it. -
max_entries- Maximum total entries to retain. Oldest are trimmed first.Every entry also records the per-file validation results
stale,changedandnew_history_rowsas set on the :class:ObjectInfo(nullwhen the run did not determine them).
ComparisonSummary Objects
@dataclass(frozen=True)
class ComparisonSummary()
Outcome of comparing a file list with the tracker, as facts to log or act on.
unchanged, changed, new and not_comparable hold the stable
keys of the current files; gone holds keys the tracker knows but the
source no longer delivers. by_md5, by_size and by_mtime count on
which basis a verdict was reached.
basis
@property
def basis() -> str
What the verdicts actually rest on: every basis that was reached, in
md5, size, mtime order and joined with + when a batch
took more than one route, or "none" when nothing could be compared.
Sources that upload without a Content-MD5 header leave both sides
with sizes only, so the conclusion is weaker than an MD5 comparison —
hence reporting the basis rather than assuming it. Every basis
:func:_compare_with_previous can return is named here; a batch decided
on modification time reads mtime, not none.
compare_with_tracker
def compare_with_tracker(file_list: list[ObjectInfo],
previous: list[dict]) -> ComparisonSummary
Compares the current file list with the latest tracker entry per file and
returns the outcome, setting changed on every :class:ObjectInfo along
the way.
Reaches its verdict without logging it, so the caller decides what to report and at which level.
Arguments:
file_list- Current list of :class:ObjectInfofrompull_files.previous- All entries loaded by :func:load_previous_snapshot.
Returns:
ComparisonSummary- Per-category keys and the basis counts.
format_names
def format_names(names: list[str], limit: int = 10) -> str
Renders keys for a log line, capped at limit names with a remainder
count so a batch of hundreds stays readable.
hours_since_last_load
def hours_since_last_load(previous: list[dict]) -> float | None
Hours since the newest tracker entry, or None when the tracker holds no
usable timestamp.
Entries are written only by runs that actually loaded — a skipped run
returns before the snapshot is saved — so the newest captured_at marks
the last full load rather than the last run.
Arguments:
previous- All entries loaded by :func:load_previous_snapshot.
Returns:
float | None: Age of the most recent load in hours.
log_comparison_summary
def log_comparison_summary(summary: ComparisonSummary) -> None
Reports a :class:ComparisonSummary: always the counts, and the file names
only when something deviates.
Kept apart from :func:compare_with_tracker so both bronze paths — the
pre-check and the post-download check — report the same way while the
comparison itself stays free of logging.