id: ml-online-feature-store-sync
namespace: company.team
description: |
Turn a Kafka stream of feature events into a serving-ready online store in
Redis and a queryable offline snapshot in DuckDB/parquet, with schema-drift
and freshness gates so models never train or serve on silently broken
features.
triggers:
- id: on_feature_event
type: io.kestra.plugin.kafka.RealtimeTrigger
description: Consume each feature event as it arrives and start one execution
per message. Enable once your broker and topic are set.
topic: feature-updates
properties:
bootstrap.servers: "{{ secret('KAFKA_BOOTSTRAP_SERVERS') }}"
groupId: kestra-feature-store
valueDeserializer: JSON
disabled: true
inputs:
- id: required_features
type: STRING
displayName: Required features
description: Comma-separated feature names every event must carry as numbers
before anything is written to the serving store.
defaults: "clicks,purchases,amount_sum"
- id: window_seconds
type: INT
displayName: Window seconds
description: Tumbling window size used to bucket events - event times are
floored to this multiple of the epoch.
defaults: 600
- id: max_lag_seconds
type: INT
displayName: Max lag seconds
description: Events older than this are flagged stale and page the channel
instead of poisoning the online store.
defaults: 900
tasks:
- id: compute_window_features
type: io.kestra.plugin.scripts.python.Script
description: Validate the event against the required feature contract, bucket it
into the tumbling window, derive serving features, and measure freshness
lag.
taskRunner:
type: io.kestra.plugin.core.runner.Process
env:
EVENT_JSON: "{{ trigger.value | toJson }}"
REQUIRED_FEATURES: "{{ inputs.required_features }}"
WINDOW_SECONDS: "{{ inputs.window_seconds }}"
MAX_LAG_SECONDS: "{{ inputs.max_lag_seconds }}"
script: |
import json
import os
import time
from datetime import datetime, timezone
def emit(outputs):
print("::" + json.dumps({"outputs": outputs}) + "::")
def parse_int(name, default):
try:
value = int(float(os.environ.get(name, "") or default))
return value if value > 0 else default
except Exception:
return default
def now_epoch():
raw = (os.environ.get("NOW_EPOCH", "") or "").strip()
if raw:
try:
return float(raw)
except Exception:
pass
return time.time()
def parse_time(raw):
if isinstance(raw, (int, float)) and not isinstance(raw, bool):
return float(raw)
text = str(raw or "").strip()
if not text:
return None
try:
return float(text)
except Exception:
pass
try:
if text.endswith("Z"):
text = text[:-1] + "+00:00"
return datetime.fromisoformat(text).timestamp()
except Exception:
return None
def iso(epoch):
return datetime.fromtimestamp(epoch, tz=timezone.utc).isoformat()
window_seconds = parse_int("WINDOW_SECONDS", 600)
max_lag_seconds = parse_int("MAX_LAG_SECONDS", 900)
required = [f.strip() for f in (os.environ.get("REQUIRED_FEATURES", "") or "").split(",") if f.strip()]
raw = (os.environ.get("EVENT_JSON", "") or "").strip()
malformed = False
event = None
if not raw:
malformed = True
else:
try:
parsed = json.loads(raw)
event = parsed if isinstance(parsed, dict) else None
except Exception:
event = None
malformed = event is None
drift = []
features = {}
rows = []
entity = ""
bucket_start = ""
bucket_end = ""
lag_seconds = 0
stale = False
if malformed:
drift = ["malformed_event"]
else:
entity = str(event.get("entity_id") or event.get("id") or "").strip()
if not entity:
drift.append("entity_id")
window = event.get("window") if isinstance(event.get("window"), dict) else {}
event_epoch = parse_time(event.get("event_time"))
if event_epoch is None:
drift.append("event_time")
event_epoch = now_epoch()
bucket = (int(event_epoch) // window_seconds) * window_seconds
bucket_start = iso(bucket)
bucket_end = iso(bucket + window_seconds)
for name in required:
if name not in window:
drift.append(name)
continue
value = window[name]
if isinstance(value, bool) or not isinstance(value, (int, float)):
drift.append(name)
continue
features[name] = float(value)
purchases = features.get("purchases")
clicks = features.get("clicks")
amount_sum = features.get("amount_sum")
if purchases is not None and amount_sum is not None:
features["amount_per_purchase"] = round(amount_sum / purchases, 4) if purchases > 0 else 0.0
if purchases is not None and clicks is not None:
features["conversion_rate"] = round(purchases / clicks, 6) if clicks > 0 else 0.0
lag_seconds = max(0, int(now_epoch() - event_epoch))
stale = lag_seconds > max_lag_seconds
if features and not drift:
rows = [{
"entity_key": entity,
"bucket_start": bucket_start,
"bucket_end": bucket_end,
"lag_seconds": lag_seconds,
"features": features,
}]
with open("features.ndjson", "w", encoding="utf-8") as out:
for row in rows:
out.write(json.dumps(row) + "\n")
emit({
"entity_key": entity or "unknown",
"features": features,
"drift_fields": drift,
"stale": stale,
"lag_seconds": lag_seconds,
"bucket_start": bucket_start,
"bucket_end": bucket_end,
"feature_count": len(features),
"malformed": malformed,
})
- id: route_schema
type: io.kestra.plugin.core.flow.If
description: Missing, non-numeric, or malformed events page the channel and skip
both stores - a broken contract must never reach serving.
condition: "{{ outputs.compute_window_features.vars.drift_fields | length > 0 }}"
then:
- id: slack_drift
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Report which features broke the contract for this event.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Feature store schema drift on {{ trigger.topic }}: entity {{ outputs.compute_window_features.vars.entity_key }} failed on {{ outputs.compute_window_features.vars.drift_fields | join(', ') }}. Online and offline writes skipped for this event (execution {{ execution.id }})."
}
else:
- id: write_online
type: io.kestra.plugin.redis.string.Set
description: Publish the derived feature vector so online models can serve from
Redis.
url: "{{ secret('REDIS_URL') }}"
key: "feat:{{ outputs.compute_window_features.vars.entity_key }}"
value: "{{ outputs.compute_window_features.vars.features | toJson }}"
- id: export_offline
type: io.kestra.plugin.jdbc.duckdb.Query
description: Archive this window as a parquet snapshot for training and backfills.
inputFiles:
features.ndjson: "{{
outputs.compute_window_features.outputFiles['features.ndjson'] }}"
outputFiles:
- features.parquet
sql: |
COPY (
SELECT * FROM read_json_auto('features.ndjson')
) TO 'features.parquet' (FORMAT PARQUET)
- id: route_freshness
type: io.kestra.plugin.core.flow.If
description: Events older than the freshness budget page the channel - training
on a stale window is worse than skipping it.
condition: "{{ outputs.compute_window_features.vars.stale == true }}"
then:
- id: slack_stale
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Flag the freshness breach with the measured lag.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Feature store freshness breach: entity {{ outputs.compute_window_features.vars.entity_key }} arrived {{ outputs.compute_window_features.vars.lag_seconds }}s late (max {{ inputs.max_lag_seconds }}s) for window {{ outputs.compute_window_features.vars.bucket_start }} (execution {{ execution.id }})."
}
- id: log_summary
type: io.kestra.plugin.core.log.Log
description: One line per event so the execution history doubles as an ingest log.
message: "feature_store entity={{
outputs.compute_window_features.vars.entity_key }} window={{
outputs.compute_window_features.vars.bucket_start }} features={{
outputs.compute_window_features.vars.feature_count }} drift={{
outputs.compute_window_features.vars.drift_fields | length }} stale={{
outputs.compute_window_features.vars.stale }} lag={{
outputs.compute_window_features.vars.lag_seconds }}s"
errors:
- id: slack_errors
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Page the channel when the flow itself fails - a silent consumer is
a silent training-set hole.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "ml-online-feature-store-sync FAILED on {{ trigger.topic | default('unknown source') }} (execution {{ execution.id }}): {{ flowRevision is defined and error is defined ? error : 'see execution logs' }}"
}
outputs:
- id: entity_key
type: STRING
value: "{{ outputs.compute_window_features.vars.entity_key }}"
- id: features_written
type: BOOL
value: "{{ outputs.write_online is defined ? true : false }}"
- id: drift_fields
type: JSON
value: "{{ outputs.compute_window_features.vars.drift_fields }}"
- id: stale
type: BOOL
value: "{{ outputs.compute_window_features.vars.stale }}"
- id: lag_seconds
type: INT
value: "{{ outputs.compute_window_features.vars.lag_seconds }}"
- id: bucket_start
type: STRING
value: "{{ outputs.compute_window_features.vars.bucket_start }}"