Skip to main content

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:ObjectInfo list from pull_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, changed and new_history_rows as set on the :class:ObjectInfo (null when 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 &quot;none&quot; 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:ObjectInfo from pull_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.