Dataset Onboard
Bring a new raw dataset into the platform with schema documentation, profiling, and an ingestion job.
Inputs
- Dataset name (snake_case, e.g.,
customer_events) - Source format (
csv,parquet,json,avro) - Source path or URI (S3, GCS, local mount, API endpoint)
- Expected frequency (
daily,weekly,on-demand) - Owner team and contact
Steps
Schema sniff
import pandas as pd df = pd.read_csv("{source_path}", nrows=1000) # or read_parquet, etc. print(df.dtypes) print(df.describe(include="all")) print(df.isnull().sum() / len(df)) # null ratesDocument:
- Column names, inferred types, null rates, example values
- Detected anomalies (mixed types, encoding issues, unexpected nulls)
Create the data dictionary entry Edit
docs/data-dictionary/{dataset_name}.md:# {DatasetName} **Owner:** {team} **Contact:** {email} **Source:** {uri} **Frequency:** {frequency} | Column | Type | Nullable | Description | |--------|------|----------|-------------| | ... | ... | ... | ... |Generate the profiling notebook
cp templates/profiling-notebook.ipynb \ notebooks/profiling/{dataset_name}_profile.ipynbEdit the notebook to use the correct source path and column list. Run it to confirm it completes without errors.
Create the ingestion job
pipelines/ingestion/{dataset_name}/ ├── ingest.py # main ingestion script ├── schema.py # column definitions and type coercion ├── config.yaml # source path, schedule, destination table └── tests/ └── test_ingest.py # unit test with a small fixture fileThe ingestion script must be idempotent (re-running on the same input produces the same output; no duplicate rows).
Register in the scheduler Add a DAG entry in
dags/{dataset_name}_ingest.py(Airflow) or a flow inflows/(Prefect). Set the schedule to match{frequency}.Test the ingestion
python -m pytest pipelines/ingestion/{dataset_name}/tests/ -v python pipelines/ingestion/{dataset_name}/ingest.py --dry-runRun a data quality gate (invoke
data-quality-gateskill) Add at minimum: null checks for required columns, row-count sanity check.
Conventions
- Raw datasets land in the
raw/schema/layer; never write tostaging/ormarts/from an ingestion job - All ingestion jobs accept
--dry-runand--dateflags - Dataset names are snake_case; tables follow the same naming
- Profiling notebooks live in
notebooks/profiling/; they are committed
Edge Cases
- Malformed source file: Log the error with row number, skip the bad row, emit a
data_quality_alertmetric. Never silently swallow rows. - Schema drift (columns added/removed): The ingestion job must compare the inbound schema to
schema.pyand fail fast on unexpected changes, rather than silently loading partial data. - Large datasets (>1GB): Use chunked reads (
chunksizein pandas, or native Parquet partitioning); test with a 10k-row sample first. - API source with rate limits: Add exponential back-off and a
--resume-fromflag that uses a checkpoint file.