loaders.pull_files
pull_files
def pull_files(table_config: TableConfig,
config_manager: ConfigManager = None) -> list[ObjectInfo]
Pulls files from a remote source to a local destination based on the configuration and connection type provided. It handles various connection types, retrieves the file details, processes them, and returns a list of ObjectInfo instances representing the pulled files.
The connection type decides what happens, and it is read from
table_config.connectiontype when set, otherwise from the connection the object
points at:
skipbronze— nothing is pulled; the bronze folder is assumed to be filled by something else and an empty list is returned.fabricfiles— the files already sit in the bronze lakehouseA place where you store both "raw" data (like files) and "organized" data (like tables). It combines the best of a File Cabinet and a Database. and are only listed, not copied. Sub-directories (such as_tracking/) and*_fileinfo.jsonsidecars are skipped, and a source path that does not exist yet yields an empty list instead of an error.azblob— the files are downloaded from Azure Blob Storage, see below.customperfile— requiresconnectiontypeon the object itself; without it the call fails rather than guessing.
Azure Blob Storage path resolution (azblob):
- The connection string comes from Key Vault, never from configuration: the secret
named by
connection.keyvaultsecretconnectionstringis read fromhttps://{config_manager.keyvault}.vault.azure.net/. - The source side is
connection.containerplus the source folder —table_config.sourcefolderwhen set, otherwiseconnection.sourcefolder. Only that folder is listed, never its sub-folders. Which blobs match is decided bytable_config.sourcefilter, a regular expression that defaults to^{sourcetable}$, andtable_config.sourceorderfixes the order they are handled in. - Each downloaded blob is renamed to
{HHMM}_{random}_{basename}— the timestamp and random suffix keep two runs on the same day from overwriting each other. - The destination is derived, not configured. Both paths get the same shape, the
abfs one from the
abfspathof the bronze layer and the mount one from itsmountpath:{bronze path}/Files/{bronzefolder}/{dataplatformobjectname}/{yyyymmdd}/{bronzefilename}, wherebronzefolderagain falls back fromtable_configto the connection, andyyyymmddis the run date. So the per-day folder is named after the data platform object, not after the source table or the source folder. - An empty match is not an error: a single
ObjectInfonamedno_filesis returned, carrying the filter and container in its log so the run shows why nothing arrived.
Arguments:
config_manager- Manages configurations used throughout the application.table_config- Contains specific information and settings for a particular table including connection type, source folder, and bronze folder.
Returns:
List[ObjectInfo]- A list containing ObjectInfo instances, each representing a pulled file with relevant details such as its bronze file name, path, and source log.
Raises:
Exception- If a connection type is not set on the object for "customperfile".Exception- If an unsupported connection type is detected.
precheck_source
def precheck_source(
table_config: TableConfig,
previous: list[dict],
config_manager: ConfigManager = None,
violation_policy=None,
skip_stale_check: bool = False) -> ComparisonSummary | None
Pre-checks Azure Blob source files against the previous tracker snapshot, WITHOUT downloading any files.
Only applicable for the azblob connection type. Returns None — no
comparison was made, so the caller must load — for any other connection
type, when previous is empty (first run), and when the source filter
matches no files at all.
Uses :func:list_blobs to retrieve current blob metadata and compares it
against previous via :func:~.file_tracker.compare_with_tracker. The
basename of each blob name is the stable key, aligning with the
partial_filename stored in the tracker.
Returning the summary rather than a verdict lets the caller report what the
conclusion rests on: a source that uploads without a Content-MD5 header
is compared on size alone.
When the source is unchanged the listed blob metadata still carries each
blob's last_modified, so the freshness check runs here before the load
is skipped — otherwise a supplier that stopped delivering (old file, never
replaced) would be skipped silently. The check is gated by
skip_stale_check and routed through violation_policy so it warns or
stops per the configured action; pass neither to skip freshness entirely.
Arguments:
config_manager- Application-level ConfigManager.table_config- TableConfig for the table being processed.previous- Tracker entries from :func:~.file_tracker.load_previous_snapshot.violation_policy- Policy used to report a stale-but-unchanged source.skip_stale_check- WhenTrue, skip the freshness check.
Returns:
ComparisonSummary | None: The comparison outcome, or None when no
comparison could be made and the caller must load.
Raises:
FabricStaleFileError- When the source is unchanged but older thanbronzemaxfileagehoursand the violation policy stops on error.
age_in_hours
def age_in_hours(moment: datetime) -> float
Hours between moment and now, both compared in UTC.
report_stale_files
def report_stale_files(files: list[ObjectInfo], table_config: TableConfig,
violation_policy) -> None
Flag every file older than bronzemaxfileagehours and report it (Validation 1).
Two thresholds, one verdict per file: past bronzemaxfileagehours at the base
severity, past bronzemaxfileagecriticalhours as critical instead.
Escalation moves the severity, never the action.
Called by both the blob pre-check and the post-download check.
file_modified_at
def file_modified_at(fileinfo) -> datetime | None
Modification time of a notebookutils FileInfo as timezone-aware UTC:
modifyTime (ms since epoch), or modificationTime on older builds.