惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Threat Intelligence Blog | Flashpoint
Threat Intelligence Blog | Flashpoint
腾讯CDC
T
Threatpost
L
Lohrmann on Cybersecurity
P
Proofpoint News Feed
The Cloudflare Blog
博客园 - 聂微东
MyScale Blog
MyScale Blog
M
MIT News - Artificial intelligence
Hacker News - Newest:
Hacker News - Newest: "LLM"
S
SegmentFault 最新的问题
博客园 - 三生石上(FineUI控件)
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
Project Zero
Project Zero
Simon Willison's Weblog
Simon Willison's Weblog
S
Schneier on Security
Spread Privacy
Spread Privacy
The GitHub Blog
The GitHub Blog
D
DataBreaches.Net
S
Securelist
Schneier on Security
Schneier on Security
Microsoft Azure Blog
Microsoft Azure Blog
cs.CL updates on arXiv.org
cs.CL updates on arXiv.org
Cisco Talos Blog
Cisco Talos Blog
博客园 - 叶小钗
量子位
I
InfoQ
J
Java Code Geeks
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
人人都是产品经理
人人都是产品经理
博客园_首页
C
CERT Recently Published Vulnerability Notes
I
Intezer
Y
Y Combinator Blog
T
Tailwind CSS Blog
Microsoft Security Blog
Microsoft Security Blog
Application and Cybersecurity Blog
Application and Cybersecurity Blog
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
大猫的无限游戏
大猫的无限游戏
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
A
About on SuperTechFans
A
Arctic Wolf
阮一峰的网络日志
阮一峰的网络日志
P
Proofpoint News Feed
G
Google Developers Blog
C
Cyber Attacks, Cyber Crime and Cyber Security
Vercel News
Vercel News
PCI Perspectives
PCI Perspectives
The Last Watchdog
The Last Watchdog
月光博客
月光博客

DEV Community

Authentication Security Deep Dive: From Brute Force to Salted Hashing (With Java Examples) Why AI Systems Don’t Fail — They Drift Spilling beans for how i learn for exam😁"Reinforcement Learning Cheat Sheet" I Replaced Chrome with Safari for AI Browser Automation. Here's What Broke (and What Finally Worked) How Python Borrows Other People's Work The $40 Architecture: Processing 1 Billion API Requests with 99.99% Uptime Vibe Coding: A Workflow Guide (From Zero to SaaS) Most webhook security guides protect the wrong side. The scary part is delivery. Headless CMS for TanStack Start: Build a Blog with Cosmic EU Age Verification App "Hacked in 2 Minutes" — What Actually Happened Comfy Cloud’s delete function does not actually remove files Running AI Models on GPU Cloud Servers: A Beginner Guide Event-driven media intelligence with AWS Step Functions and Bedrock I scored 500 AI prompts across 8 quality dimensions — here's what broke How to Call Google Gemini API from Next.js (Free Tier, No Backend Needed) The Portal Protocol: Reclaiming Human Connection in the Age of AI How to Fix Your Team's Scattered Knowledge Problem With a Self-Hosted Forum Intro to tc Cloud Functors: A Graph-First Mental Model for the Modern Cloud Designing Multi-Tenant Backends With Both Ownership and Team Access I Built a Neumorphic CSS Library with 77+ Components — Here's What I Learned PostgreSQL Performance Optimization: Why Connection Pooling Is Critical at Scale Cómo construí un SaaS multi-rubro para gestionar expensas en Argentina con FastAPI + Vue 3 🚀 I Built an Ethical Hacking Scanner Tool – Open Source Project I Replaced /usage and /context in Claude Code With a Single Statusline A Pythonic Way to Handle Emails (IMAP/SMTP) with Auto-Discovery and AI-Ready Design I Collected 8.9 Million Polymarket Price Points — Here's What I Found About How Markets Really Move EcoTrack AI — Carbon Footprint Tracker & Dashboard Everyone's Using AI. No One Agrees How. 5 self-hosted ebook managers worth trying in 2026 Building Your First AI Agent with LangChain: From Chatbot to Autonomous Assistant Common SOC 2 Failures (Real World) Stop Vibe-Checking Your AI App: A Practical Guide to Evals How to Use SonarQube and SonarScanner Locally to Level Up Your Code Quality Your Next To-Do App Is Dead — I Replaced Mine with an OpenClaw AI Sign a Nostr event in 60 lines of Python using coincurve — no nostr-sdk, no nbxplorer, no rust toolchain ITGC Audit Explained Like You’re in Big 4 Patch Tuesday abril 2026: Microsoft parcha 163 vulnerabilidades y un zero-day en SharePoint Stop scraping everything: a better way to track competitor price changes Listing on MCPize + the Official MCP Registry while routing payments OUTSIDE the marketplace — how I kept 100% of my x402 revenue Building an AI-Powered Risk Intelligence System Using Serverless Architecture Why We Ripped Function Overloading Out of Our AI Toolchain Testing AI-Generated Code: How to Actually Know If It Works SaaS Churn Is Killing Your Business. Here Is What to Do About It (Without a Support Team) The Speed of AI Is No Longer Linear - And Self-Improving Models Are Why How to Implement RBAC for MCP Tools: A Practical Guide for Engineering Teams From Standard Quote to Persuasive Proposal: AI Automation for Arborists I built a CLI that scaffolds complete multi-tenant SaaS apps Axios CVE-2025–62718: The Silent SSRF Bug That Could Be Hiding in Your Node.js App Right Now The dashboard that ended our friendship Data Pipelines Explained Simply (and How to Build Them with Python) The Hidden Cost of AI Systems Nobody Talks About. undefined vs undeclared, and how typeof behaves Switching from file-based jobs to NATS/Kafka in Rust without changing code io_uring Adventures: Rust Servers That Love Syscalls Why Agentic AI is Killing the Traditional Database The POUR principles of web accessibility for developers and designers Quantum Neural Network 3D — A Deep Dive into Interactive WebGL Visualization How To Install Caveman In Codex On macOS And Windows Automation Pipeline Reliability: Why Your Workflow Breaks When Nobody Is Watching I Built an 'Open World' AI Coding Agent — It Works From ANY Folder From Freelancing to Product: A Tech Service Company's SaaS Transformation China's AI Giants: Adding Tencent Hunyuan & ByteDance Doubao to AI University (74 Providers) On the Vibe Coders and Their Lies clerk: Auto-Summarize Your Claude Code Sessions AI Weekly — 2026/04/10–04/17 | The Model Lockdown Is Here, but the Toolchain Is the Real Battleground AI 週報 — 2026/04/10–2026/04/17 模型封鎖潮來了,但工具鏈才是真戰場 Maybe this is how Open-Source apps are born... 🚀 Fine-Tune LLMs with LoRA and QLoRA: 2026 Guide tRPC v11 + Next.js App Router: End-to-End Type Safety Without the Boilerplate ShadCN UI in 2026: Why I Stopped Installing Component Libraries and Started Owning My Components SaaS Billing in React Server Components: Stripe + Supabase Without a Single `useEffect` Join our DEV Weekend Challenge — $1,000 in Prizes Across TEN winners! Submissions Due April 20 at 6:59 AM UTC. Implementing FSRS Spaced Repetition in Flutter + Supabase — Adding Memory Science to an AI Learning App "I Texted My Localhost From the Train — Claude Code Fixed the Bug Before I Got Home" I Built a Sales Prep AI and It Went Deeper Than Expected Design to Code #2: One JSON, Eleven Outputs Solving the 100M-Row Problem: A Summary Table Pattern for High-Volume Push Notification Logs Flutter Web With Wasm: What Actually Changes For Developers I Built 50 Royalty-Free Soundtracks for My Side Project in a Weekend Using AI Music Generation The Vibe Coding Security Checklist: 7 Things to Check Before You Ship Stop Letting Googlebot Guess Fix Your React App's SEO Right Desconstruindo o Streaming do LinkedIn: Como Criar um Engine de Extração de Vídeo de Alta Performance com HLS e FFmpeg (EDA Part-1) EDA (Exploratory Data Analysis) Explained With Real Life — Why Looking at Your Data Is the Most Important Step in Machine Learning Brand Relationship Management at Scale: Our 4-Touch Outreach System for 200+ Brands Why String.fromEnvironment() Might Return an Empty String in Dart JGuardrails 1.0.0 — Hardening Java LLM Apps Against Jailbreaks, Toxicity, and Prompt Injection Plan and Schedule a Full Week of Threads Content From One Claude Conversation Coding Cat Oran Ep3, Five Tables Changed Everything Updated: BFF Pattern I'm done watching freelancers get buried by 200 proposals. So I'm building the alternative. This is my first post BFS Algorithm in Java Step by Step Tutorial with Examples Tracking LLM Pricing Monthly: An Open Dataset for 22 AI Models How We Measure Content ROI on a Comparison Site: Revenue Attribution Without Perfect Data Introducing Nova AI Ops: The AI-Native Operating System for SRE Teams I built a free desktop video downloader for Windows — Grabbit How Talkie OCR Helps Vision-Impaired & Dyslexic Users Read the World Around Them VRCFaceTracking安装和iPhone面捕配置教程,有bug Even CrowdStrike Can't See Your Agents The Automation Gold Rush: What n8n Workflows and Claude Are Opening Up for Developers Right Now
Real-Time AI Feature Engineering with Spark Structured Streaming and Databricks Feature Store
Jubin Soni · 2026-06-24 · via DEV Community

Building point-in-time correct, production-grade feature pipelines — from raw Kafka events to online feature serving in milliseconds, using Spark Structured Streaming and the Databricks Feature Store.


Table of Contents

  • The Feature Engineering Problem
  • Architecture Overview
  • Feature Store Concepts: ERD
  • Environment Setup
  • Streaming Feature Pipeline
  • Point-in-Time Correct Training Dataset Generation
  • Writing Features to the Online Store
  • Serving Features at Inference Time
  • Feature Table Reference
  • References

The Feature Engineering Problem

Feature engineering is where most ML projects silently fail in production. Not because the model is wrong — but because the features the model sees at training time are different from the features it sees at inference time. This is called training-serving skew, and it's the #1 silent killer of ML systems.

Three specific failure modes cause it:

  1. Online/offline inconsistency — the batch pipeline that computes training features uses different logic than the real-time service that computes inference features
  2. Data leakage — training features accidentally include information from the future (e.g. joining on a label that was created after the event)
  3. Feature staleness — a model trained on 30-day rolling averages is served features that are 6 hours stale because the pipeline backfills are slow

The Databricks Feature Store — now part of Unity Catalog as Feature Engineering in Unity Catalog — solves all three by:

  • Storing feature computation logic alongside the data (no drift between training and serving)
  • Enforcing point-in-time lookups during training dataset creation
  • Providing a unified API for both batch offline reads and low-latency online reads

Architecture Overview

Architecture description

Feature Store Concepts: ERD

Understanding the data model behind the Feature Store is essential for designing correct pipelines. Here's how the entities relate:

ERD description

The critical relationship: a Model Version is bound to a Training Set, which records exactly which feature tables and which point-in-time lookups were used. This is how Databricks guarantees reproducibility — you can always re-create the exact training data that produced any model version.


Environment Setup

# Databricks Runtime ML 13.x+ recommended
# Feature Engineering in Unity Catalog (formerly Feature Store)

%pip install databricks-feature-engineering==0.6.0 --quiet
dbutils.library.restartPython()

from databricks.feature_engineering import FeatureEngineeringClient, FeatureLookup
from databricks.feature_engineering.entities.feature_serving_endpoint import (
    ServedEntity, EndpointCoreConfig
)
from pyspark.sql import functions as F, SparkSession
from pyspark.sql.types import (
    StructType, StructField, StringType, LongType,
    DoubleType, TimestampType, ArrayType
)
import mlflow

spark = SparkSession.builder.getOrCreate()
fe = FeatureEngineeringClient()

# Unity Catalog paths
CATALOG       = "prod"
FEATURE_DB    = f"{CATALOG}.feature_store"
EVENTS_TABLE  = f"{CATALOG}.silver.events_clean"
KAFKA_BROKER  = "kafka-broker.internal:9092"
KAFKA_TOPIC   = "user-events"

# Checkpoint locations (ADLS / S3 / GCS)
CHECKPOINT_BASE = "abfss://checkpoints@storage.dfs.core.windows.net/features"


Streaming Feature Pipeline

The streaming pipeline reads from Kafka, computes windowed aggregations using Spark's stateful streaming engine, and writes features to the Feature Store via foreachBatch. This keeps the feature table continuously fresh.

# ── Streaming Feature Pipeline ────────────────────────────────────────────────

# Step 1: Define the raw event schema from Kafka
event_schema = StructType([
    StructField("user_id",       StringType(),    False),
    StructField("event_type",    StringType(),    True),
    StructField("product_id",    StringType(),    True),
    StructField("revenue",       DoubleType(),    True),
    StructField("session_id",    StringType(),    True),
    StructField("platform",      StringType(),    True),
    StructField("event_ts",      TimestampType(), False),
])

# Step 2: Read from Kafka
raw_stream = (
    spark.readStream
        .format("kafka")
        .option("kafka.bootstrap.servers", KAFKA_BROKER)
        .option("subscribe", KAFKA_TOPIC)
        .option("startingOffsets", "latest")
        .option("failOnDataLoss", "false")
        .load()
        .select(
            F.from_json(F.col("value").cast("string"), event_schema).alias("data"),
            F.col("timestamp").alias("kafka_ts")
        )
        .select("data.*", "kafka_ts")
)

# Step 3: Apply watermark and compute windowed features
# Watermark: tolerate up to 10 minutes of late data
windowed_features = (
    raw_stream
        .withWatermark("event_ts", "10 minutes")
        .groupBy(
            F.col("user_id"),
            F.window(F.col("event_ts"), "1 hour", "15 minutes").alias("window")
        )
        .agg(
            F.count("*").alias("event_count_1h"),
            F.sum(F.when(F.col("event_type") == "purchase", F.col("revenue"))
                  .otherwise(0)).alias("revenue_1h"),
            F.countDistinct("session_id").alias("session_count_1h"),
            F.countDistinct("product_id").alias("unique_products_1h"),
            F.sum(F.when(F.col("event_type") == "purchase", 1)
                  .otherwise(0)).alias("purchase_count_1h"),
            F.first("platform").alias("last_platform"),
        )
        # Flatten window struct to scalar columns
        .withColumn("window_start", F.col("window.start"))
        .withColumn("window_end",   F.col("window.end"))
        .withColumn("feature_ts",   F.col("window.end"))   # timestamp key for PIT lookup
        .drop("window")
        # Derived features
        .withColumn("conversion_rate_1h",
            F.when(F.col("event_count_1h") > 0,
                   F.col("purchase_count_1h") / F.col("event_count_1h"))
            .otherwise(0.0))
        .withColumn("avg_revenue_per_purchase_1h",
            F.when(F.col("purchase_count_1h") > 0,
                   F.col("revenue_1h") / F.col("purchase_count_1h"))
            .otherwise(0.0))
)


# Step 4: Write to Feature Store via foreachBatch
# foreachBatch gives us transactional writes per micro-batch
def write_to_feature_store(batch_df, batch_id):
    """
    Called on each micro-batch. Merges feature data into the Feature Store
    table using merge_on keys (user_id + feature_ts).
    """
    if batch_df.isEmpty():
        return

    fe.write_table(
        name=f"{FEATURE_DB}.user_activity_features",
        df=batch_df,
        mode="merge",             # upsert: update existing, insert new
    )
    print(f"Batch {batch_id}: wrote {batch_df.count()} feature rows")


# Step 5: Create the feature table (idempotent — safe to re-run)
try:
    fe.create_table(
        name=f"{FEATURE_DB}.user_activity_features",
        primary_keys=["user_id"],
        timestamp_keys=["feature_ts"],
        schema=windowed_features.schema,
        description=(
            "Real-time user activity features computed from event stream. "
            "1-hour sliding window, refreshed every 15 minutes. "
            "Primary key: user_id. Timestamp key: feature_ts (window end)."
        ),
    )
    print("Feature table created.")
except Exception:
    print("Feature table already exists — continuing.")


# Step 6: Launch the streaming query
streaming_query = (
    windowed_features.writeStream
        .outputMode("update")               # update mode for stateful aggregations
        .option("checkpointLocation", f"{CHECKPOINT_BASE}/user_activity")
        .trigger(processingTime="5 minutes") # micro-batch every 5 min
        .foreachBatch(write_to_feature_store)
        .start()
)

print(f"Streaming query '{streaming_query.name}' running...")
print(f"Status: {streaming_query.status}")


Point-in-Time Correct Training Dataset Generation

This is the most critical part of the Feature Store. When creating training data, we must join labels to features at the timestamp of the label event — not the current time. This prevents data leakage.

# ── Point-in-Time Correct Training Dataset ────────────────────────────────────

# Step 1: Load the label dataset
# Each row = one prediction target event, with the exact timestamp
# at which a model would have needed to make a prediction.

labels_df = (
    spark.table(f"{CATALOG}.gold.churn_labels")
        .select(
            "user_id",
            "churn_label",                        # 0 = retained, 1 = churned
            F.col("observation_ts").alias("event_timestamp"),  # point-in-time anchor
            "experiment_split"                    # train/val/test
        )
        .filter(F.col("observation_ts") >= "2024-01-01")
)

print(f"Label rows: {labels_df.count():,}")
labels_df.show(5)
# +----------+-----------+---------------------+-----------------+
# | user_id  |churn_label| event_timestamp     | experiment_split|
# +----------+-----------+---------------------+-----------------+
# | u_123456 | 0         | 2024-03-15 14:22:00 | train           |
# | u_789012 | 1         | 2024-03-15 18:45:00 | train           |


# Step 2: Define feature lookups
# as_of_timestamp=None → use the label's event_timestamp (point-in-time)
# Databricks will join each label row to the feature values
# that were valid at event_timestamp — not the latest values.

feature_lookups = [
    # User activity features — 1h window features from the streaming pipeline
    FeatureLookup(
        table_name=f"{FEATURE_DB}.user_activity_features",
        feature_names=[
            "event_count_1h",
            "revenue_1h",
            "session_count_1h",
            "unique_products_1h",
            "purchase_count_1h",
            "conversion_rate_1h",
            "avg_revenue_per_purchase_1h",
            "last_platform",
        ],
        lookup_key="user_id",
        timestamp_lookup_key="event_timestamp",    # ← PIT anchor
    ),

    # User profile features — slower-changing, from batch pipeline
    FeatureLookup(
        table_name=f"{FEATURE_DB}.user_profile_features",
        feature_names=[
            "account_age_days",
            "lifetime_revenue",
            "preferred_category",
            "subscription_tier",
        ],
        lookup_key="user_id",
        timestamp_lookup_key="event_timestamp",    # ← PIT anchor
    ),

    # Transaction aggregates — 30d and 90d rolling windows
    FeatureLookup(
        table_name=f"{FEATURE_DB}.transaction_features",
        feature_names=[
            "purchase_count_30d",
            "purchase_count_90d",
            "avg_order_value_30d",
            "days_since_last_purchase",
            "category_diversity_score",
        ],
        lookup_key="user_id",
        timestamp_lookup_key="event_timestamp",
    ),
]


# Step 3: Create training dataset (Feature Store handles the PIT join)
training_set = fe.create_training_set(
    df=labels_df,
    feature_lookups=feature_lookups,
    label="churn_label",
    exclude_columns=["observation_ts", "experiment_split"],
)

# The returned DataFrame has features + labels, PIT-correct
training_df = training_set.load_df()
print(f"Training rows: {training_df.count():,}")
print(f"Training cols: {len(training_df.columns)}")
training_df.show(3)


# Step 4: Train model and log via Feature Store (preserves lineage!)
from sklearn.ensemble import GradientBoostingClassifier
import pandas as pd

train_pdf = (
    training_df
        .filter(F.col("experiment_split") == "train")
        .drop("experiment_split", "user_id")
        .fillna(0)
        .toPandas()
)

X_train = train_pdf.drop(columns=["churn_label"])
y_train = train_pdf["churn_label"]

model = GradientBoostingClassifier(
    n_estimators=300,
    learning_rate=0.05,
    max_depth=5,
    subsample=0.8,
    random_state=42,
)

with mlflow.start_run(run_name="churn-gbm-v1") as run:
    model.fit(X_train, y_train)

    # Log model via Feature Store — this records the feature lineage
    fe.log_model(
        model=model,
        artifact_path="churn_model",
        flavor=mlflow.sklearn,
        training_set=training_set,      # ← binds model to its feature lookups
        registered_model_name=f"{CATALOG}.ml.user_churn_model",
    )
    print(f"Logged model with feature lineage. Run: {run.info.run_id}")


Writing Features to the Online Store

For real-time inference, the model needs features in milliseconds — not the seconds it takes to query Delta Lake. Databricks Feature Store can publish features to an online store (DynamoDB, Cosmos DB, MySQL, etc.) for low-latency reads.

# ── Publish Features to Online Store ─────────────────────────────────────────
# Online stores are configured per feature table.
# Here we publish user_activity_features to DynamoDB for <5ms lookups.

from databricks.feature_engineering.entities.feature_store_online_table import (
    OnlineTable, OnlineTableSpec, TriggeredSchedulingPolicy
)

# Create an online table spec (backed by a serverless real-time compute layer)
online_table_spec = OnlineTableSpec(
    primary_key_columns=["user_id"],
    source_table_full_name=f"{FEATURE_DB}.user_activity_features",
    run_triggered=OnlineTableSpec.TriggeredSchedulingPolicy(),  # sync on-demand
    # OR for continuous sync:
    # run_continuous=OnlineTableSpec.ContinuousSchedulingPolicy()
)

# Create the online table (idempotent)
online_table = fe.create_online_table(spec=online_table_spec)
print(f"Online table: {online_table.name}")
print(f"Status:       {online_table.status.detailed_state}")

# Trigger an initial sync from the offline Delta table to the online store
fe.refresh_online_table(name=f"{FEATURE_DB}.user_activity_features")


Serving Features at Inference Time

At inference time, the Feature Store SDK performs automatic feature lookups, joining the incoming request data with features from the online store before passing them to the model.

# ── Real-Time Feature Serving at Inference ────────────────────────────────────

import requests, json

WORKSPACE_URL = "https://<workspace>.azuredatabricks.net"
TOKEN = dbutils.secrets.get("prod-scope", "databricks-token")


# Option 1: Model Serving with automatic feature lookup
# When you logged the model with fe.log_model(), Databricks knows which
# features to fetch. You only send the lookup key (user_id) at inference time.

def predict_churn(user_ids: list) -> list:
    """
    Send only user_id — the serving endpoint fetches features automatically
    from the online store and runs inference.
    """
    payload = {
        "dataframe_records": [
            {"user_id": uid} for uid in user_ids
        ]
    }
    resp = requests.post(
        f"{WORKSPACE_URL}/serving-endpoints/churn-predictor/invocations",
        headers={
            "Authorization": f"Bearer {TOKEN}",
            "Content-Type":  "application/json",
        },
        data=json.dumps(payload),
        timeout=5,
    )
    resp.raise_for_status()
    return resp.json()["predictions"]


# Example usage
predictions = predict_churn(["u_123456", "u_789012", "u_345678"])
for uid, pred in zip(["u_123456", "u_789012", "u_345678"], predictions):
    print(f"{uid}: churn_probability = {pred:.4f}")
# u_123456: churn_probability = 0.0821
# u_789012: churn_probability = 0.7643
# u_345678: churn_probability = 0.1209


# Option 2: Direct feature lookup via the Feature Serving endpoint
# Useful when you want raw features without running inference
def get_features(user_ids: list) -> dict:
    payload = {
        "dataframe_records": [{"user_id": uid} for uid in user_ids]
    }
    resp = requests.post(
        f"{WORKSPACE_URL}/serving-endpoints/user-features-serving/invocations",
        headers={
            "Authorization": f"Bearer {TOKEN}",
            "Content-Type":  "application/json",
        },
        data=json.dumps(payload),
        timeout=5,
    )
    return resp.json()


# Option 3: Batch scoring (offline) — uses Delta offline store
# No online store needed; reads directly from the feature table with PIT lookup
batch_labels = spark.table(f"{CATALOG}.gold.users_to_score_today") \
    .select("user_id", F.current_timestamp().alias("event_timestamp"))

batch_predictions = fe.score_batch(
    model_uri=f"models:/{CATALOG}.ml.user_churn_model@champion",
    df=batch_labels,
    result_type="double",
)

batch_predictions.select("user_id", "prediction") \
    .write.format("delta").mode("overwrite") \
    .saveAsTable(f"{CATALOG}.gold.churn_scores_daily")


Feature Table Reference

A summary of the feature tables in our pipeline, their update cadence, and their role in the ML lifecycle:

Feature Table Primary Key Timestamp Key Update Method Latency Used In
user_activity_features user_id feature_ts Spark Structured Streaming ~5 min Real-time churn, recommendation
transaction_features user_id feature_ts Scheduled batch (hourly) ~60 min Churn, LTV prediction
user_profile_features user_id updated_at CDC from OLTP (near real-time) ~2 min All models
product_features product_id feature_ts Scheduled batch (daily) ~24 hr Recommendation, search ranking
session_features session_id session_end_ts Streaming (micro-batch) ~1 min Click-through rate, abandon prediction
cohort_features cohort_id computed_at Weekly batch ~7 days Segmentation, A/B analysis

Freshness vs cost tradeoff: Streaming features are ~10× more expensive to compute than batch features (continuous cluster vs scheduled job). Only promote a feature to streaming if your model's performance degrades meaningfully with stale data — validate this with an offline ablation study first.


Key Takeaways

  • Training-serving skew is the silent killer of production ML — the Feature Store eliminates it by encoding feature computation logic once and using it in both training and serving paths.
  • Point-in-time correct joins via timestamp_lookup_key are non-negotiable for any model trained on time-series data. A missing event_timestamp in your label table is a data leakage bug waiting to happen.
  • fe.log_model() is the right model logging call, not mlflow.sklearn.log_model(). It records feature lineage, enabling reproducible re-training and automatic feature lookup at serving time.
  • Watermarks in Structured Streaming are critical for stateful aggregations — without them, Spark accumulates state indefinitely and the job eventually OOMs. Set them to the maximum tolerable late-data window.
  • Online stores are only worth the operational cost when your SLA is under ~100ms. For batch scoring jobs or APIs with >500ms budgets, read directly from the offline Delta table.
  • fe.score_batch() is the cleanest way to run periodic batch inference — it handles PIT feature lookups automatically, keeps inference logic DRY, and logs results to Delta for downstream consumers.

References

  1. Databricks — Feature Engineering in Unity Catalog (Overview)
    🔗 https://docs.databricks.com/en/machine-learning/feature-store/uc/feature-tables-uc.html

  2. Databricks — Create and Manage Online Tables
    🔗 https://docs.databricks.com/en/machine-learning/feature-store/online-tables.html

  3. Databricks — Point-in-Time Feature Lookups
    🔗 https://docs.databricks.com/en/machine-learning/feature-store/time-series.html

  4. Apache Spark — Structured Streaming Programming Guide
    🔗 https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html

  5. Apache Spark — Streaming Watermarks for Late Data Handling
    🔗 https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html#handling-late-data-and-watermarking

  6. Databricks — Feature Store Python API Reference
    🔗 https://docs.databricks.com/en/machine-learning/feature-store/python-api.html

  7. Databricks — Score Batch with Feature Store
    🔗 https://docs.databricks.com/en/machine-learning/feature-store/score-batch.html

  8. "Feature Stores for ML" — Feast Documentation (open-source reference)
    🔗 https://docs.feast.dev/

  9. "Rethinking Feature Stores" — Chip Huyen (huyenchip.com)
    🔗 https://huyenchip.com/2023/01/08/feature-store.html

  10. Databricks — Model Serving with Automatic Feature Lookup
    🔗 https://docs.databricks.com/en/machine-learning/model-serving/feature-store-model-serving.html

  11. "Building Machine Learning Pipelines" — Hannes Hapke & Catherine Nelson (O'Reilly)
    🔗 https://www.oreilly.com/library/view/building-machine-learning/9781492053187/


This concludes the 4-part Databricks series: