Queries icon
OutputValues icon
Log icon
If icon
UploadFiles icon
Set icon
Fail icon
Schedule icon

Build a point-in-time correct training set and block feature leakage

Build ML training sets with DuckDB ASOF joins in Kestra, audit for future feature rows and correlation lift, publish only leak-free data.

Categories
AIData

The most common way to ship a model that looks great offline and fails in production is feature leakage: the training set uses a feature value from after the moment the prediction would have been made. A "join to the latest snapshot" does it silently. Every row looks fine, the model learns from the future, and nobody notices until live accuracy collapses.

This flow builds the training set point-in-time correct with a DuckDB ASOF JOIN: each label gets the newest feature value that was known at least embargo_days before its cutoff. It then audits the result before anything is published:

  • Future rows: rows whose feature is newer than the label cutoff minus the embargo. This must be 0.
  • Correlation lift: how much more the feature correlates with the label than in a point-in-time reference set. Leakage shows up as an implausibly strong signal.
  • Coverage and duplicate labels.

It runs with no setup and no network. DuckDB generates a churn dataset of 2,000 customers with weekly ticket snapshots that keep changing after each customer's cutoff, the way real feature tables do. Run it with join_strategy: LATEST to watch the gate block the leaky set: the feature/label correlation jumps from 0.23 to 0.92.

This blueprint was created by zkasuran.

How it works

  1. build_and_audit (io.kestra.plugin.jdbc.duckdb.Queries) does everything in one DuckDB session:
    • generates labels (customer, cutoff time, churned) and feature_snapshots (customer, snapshot time, open_tickets_30d),
    • builds pit_reference with ASOF LEFT JOIN ... ON f.feature_time <= l.label_time - embargo,
    • builds the candidate training set with the chosen strategy: ASOF (the reference) or LATEST (newest snapshot, the leaky pattern),
    • writes the candidate to Parquet and returns one audit row.
  2. audit (io.kestra.plugin.core.output.OutputValues) collects the numbers and computes the lift.
  3. log_audit prints the audit on every run.
  4. leakage_gate (io.kestra.plugin.core.flow.If) checks 0 future rows, 0 duplicate labels, coverage of at least min_coverage and lift of at most max_leakage_lift.
    • Pass: publish_training_set (io.kestra.plugin.core.namespace.UploadFiles) writes ml/churn/training_set.parquet, and record_manifest (io.kestra.plugin.core.kv.Set) stores what was published and its audit.
    • Fail: block_publication (io.kestra.plugin.core.execution.Fail) explains which check failed. Nothing is published.
  5. nightly (io.kestra.plugin.core.trigger.Schedule) rebuilds the set before the training job.

Inputs

  • join_strategy (SELECT ASOF or LATEST, default ASOF).
  • embargo_days (INT, default 7): minimum age of a feature at the label cutoff.
  • min_coverage (FLOAT, default 0.95): share of labels that must have a feature.
  • max_leakage_lift (FLOAT, default 0.05): allowed correlation gap to the point-in-time reference.

Quick start

Run Input Result
1 defaults Published: 2,000 labels, 0 future rows, coverage 1.0, correlation 0.2311
2 join_strategy: LATEST Blocked: 2,000 future rows, correlation 0.9218 vs 0.2311 (lift 0.69)
3 embargo_days: 70 Blocked: no leakage, but coverage 0.86. The earliest labels have no feature history that old

Expected outputs

  • outputs.audit.values: rows, leaked_rows, duplicate_labels, coverage, candidate_corr, reference_corr, leakage_lift.
  • Namespace file ml/churn/training_set.parquet.
  • KV churn_training_set_manifest: path, join, embargo, rows, coverage, correlation and execution id.

Using it with your own data

  • Replace the two generator statements with read_parquet() on your label and feature tables, or attach Postgres or a lake with the DuckDB extensions.
  • Add one ASOF LEFT JOIN per feature table. Each one keeps its own feature_time column, so the future-rows check can cover all of them.
  • Write the published set to S3, GCS or Azure instead of namespace files.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.