Data Engineer
You own: datasets, feature pipelines, and the data contracts that ML / DS / AIE consume from. GeneFlow's lineage and drift baselines depend on you.
What's new for you
| Was | Now in GeneFlow |
|---|---|
| Feature groups documented in Confluence | Declared with lineage.upstream(run, kind="feature_group", id="…") |
| Dataset versions in S3 only | Same, plus referenced by runs:/<id> lineage edges |
| Drift detection in a separate Airflow DAG | Snapshot histogram once → GeneFlow runs PSI |
| No way to find "what was trained on this?" | lineage.get_downstream(kind="feature_group", id=…) |
Day-1 setup
pip install geneflow
export GENEFLOW_TRACKING_URI=https://api.genedata.io
export GENEDATA_PAT=$(genedata auth token)
Workflow 1 — Publish a new feature group version
Your pipeline writes s3://bucket/features/user_features_v3/2026-05-12.parquet. Mark it as the new version:
import geneflow
# When a model trains on this, the run declares the upstream:
with geneflow.start_run(experiment_name="fraud_v3") as run:
geneflow.lineage.upstream(
run.run_id,
kind="feature_group",
id="user_features_v3",
)
geneflow.lineage.upstream(
run.run_id,
kind="dataset",
id="s3://bucket/features/user_features_v3/2026-05-12.parquet",
)
Workflow 2 — Compute drift baselines for downstream ML
When you publish a new feature group version, the MLE training on it needs a drift baseline. You can compute it yourself or hand them a notebook:
import numpy as np
from geneflow import serving
# Per-feature histogram
features = []
for col in ["transaction_amount", "merchant_age_days"]:
hist, edges = np.histogram(df[col], bins=10)
features.append({
"feature_name": col,
"feature_type": "numeric",
"histogram": {"buckets": hist.tolist(), "edges": edges.tolist()},
"mean": float(df[col].mean()),
"stddev": float(df[col].std()),
"min": float(df[col].min()),
"max": float(df[col].max()),
"sample_size": len(df),
})
# Save against whichever model version trained on this feature group
serving.save_drift_baselines(
model_name="fraud-detector",
model_version=3,
features=features,
)
Workflow 3 — Find downstream consumers before breaking a feature
downstream = geneflow.lineage.get_downstream(
kind="feature_group", id="user_features_v3",
)
for d in downstream:
print(d["runId"], d["modelVersion"]) # who's using it
Don't rename or remove a feature group until this list is empty (or you've coordinated migrations).
Workflow 4 — Schema contract for a feature group
GeneFlow itself doesn't enforce schemas, but you can publish them as run artifacts:
artifacts.upload(
run_id=publish_run.run_id,
src="schema.json",
dst_path="schemas/user_features_v3.json",
content_type="application/json",
)
Then downstream code can:
schema = artifacts.download(publish_run.run_id, "schemas/user_features_v3.json")
Workflow 5 — Tie a pipeline run to GeneFlow
If you run your DBT/ETL/Spark job inside an MLproject:
gfctl projects run git@github.com:genedata/feature-pipeline \
--entry-point publish-user-features-v3 \
--param date=2026-05-12 \
--backend kubernetes
The project run shows up as a gf_runs row with source_type=PROJECT. Lineage edges flow from there.
Workflow 6 — Audit "what data fed this model?"
In the UI: /dashboard/geneflow/ml/compare?run_a=… shows upstream chain on hover.
In SQL:
SELECT upstream_kind, upstream_id, produced_model_version_id
FROM gf_run_lineage
WHERE run_id = 'run_abc...'
ORDER BY recorded_at;
Common gotchas
- Lineage edges are append-only — there's no
delete_upstream()in normal SDK use; you correct mistakes via SQL with audit. feature_groupIDs are free-form strings — agree a naming convention in your org (we recommend<domain>_features_v<n>).- Don't bake row-level data into lineage IDs — they're meant for asset-level granularity (dataset, feature group, model version).
Where to go next
- 08-analytics-engineer.md — for the metric-definitions side
- 14-data-quality-engineer.md — drift baselines + DQ contracts
- api-reference.md#lineage — REST