RealtimeTrigger icon
Script icon
Process icon
If icon
SlackIncomingWebhook icon
Set icon
Query icon
Log icon

ML Online Feature Store Sync with Drift and Freshness Gates

Turn Kafka feature events into Redis online features and DuckDB offline snapshots with schema-drift and freshness gates in Kestra.

Categories
AIData

How it works

  • The on_feature_event trigger (io.kestra.plugin.kafka.RealtimeTrigger) subscribes to the feature-updates topic with a JSON deserializer and starts one execution per feature event using the kestra-feature-store consumer group.
  • compute_window_features (scripts.python.Script, Process runner) parses the event, checks every name in required_features is present and numeric, floors event_time into the window_seconds tumbling bucket, derives serving features (amount_per_purchase, conversion_rate), measures lag against max_lag_seconds, and writes one parquet-ready row to features.ndjson. Everything returns through Kestra's ::json:: outputs protocol.
  • route_schema (core.flow.If) branches on the drift list: a broken contract posts the failing feature names to Slack and skips both stores.
  • The clean path runs write_online (redis.string.Set, key feat:<entity>) so online models can serve from Redis, then export_offline (duckdb.Query) archives the window as features.parquet for training and backfills.
  • route_freshness (core.flow.If) pages the channel when the event arrived later than max_lag_seconds, and log_summary leaves one ingest line per event in the execution history.
  • Any flow-level failure lands in Slack through the errors handler.

What you get

  • A serving-ready feature vector in Redis under feat:<entity_id> on every healthy event.
  • A per-window parquet snapshot (features.parquet) archived through Kestra's internal storage.
  • Schema-drift alerts naming the exact missing or non-numeric features, freshness alerts with the measured lag, and a one-line ingest log per execution.
  • Outputs for entity_key, features_written, drift_fields, stale, lag_seconds, and bucket_start so downstream flows can gate on this one.

Who it's for

  • ML platform teams that already emit feature events to Kafka but serve features out of ad-hoc Redis writes and CSV exports, and data scientists who need a trustworthy offline copy of exactly what was served.

Why orchestrate this with Kestra

  • The windowed compute, the two writes, and the two gates are four moving parts that usually end up as a consumer service nobody owns. As a Kestra flow each event gets retries, a history, Slack alerts, and a readable If branch for the drift and freshness policies instead of buried if statements in Python.

Prerequisites

  • A Kafka broker reachable from Kestra workers and a topic (default feature-updates) carrying JSON events shaped like {"entity_id": "...", "event_time": "...", "window": {"clicks": 0, "purchases": 0, "amount_sum": 0.0}}.
  • A Redis instance reachable from Kestra workers for the online store.
  • Python on the Kestra worker (the flow uses the Process runner, so only the standard library matters).

Secrets

  • KAFKA_BOOTSTRAP_SERVERS - Kafka broker list for the trigger.
  • REDIS_URL - Redis connection URL for the online store write.
  • SLACK_WEBHOOK_URL - Incoming webhook used for drift, freshness, and error alerts.

Quick start

  • Set the three secrets above, create the topic, and set required_features to your contract (defaults cover clicks,purchases,amount_sum).
  • Enable the trigger (it ships disabled: true), publish one event, and watch the execution show the Redis key feat:<entity> plus a features.parquet output.
  • Publish an event missing a required feature and confirm the drift alert fires and both stores are skipped.

How to extend

  • Point export_offline at a persistent DuckDB file (mount a volume) to accumulate a cumulative store instead of per-window snapshots.
  • Add a second If that joins this flow's outputs with a model-registry flow (for example mlops-model-promotion-drift-watchdog) to block promotions on drift.
  • Swap the Redis write for a feature-serving API call, or loop over multiple entities per event by turning route_schema into a core.flow.Loop.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.