Проблема: training/serving skew
ML-модель в проде работает на фичах. Фичи -- это числовые и категориальные признаки, полученные из сырых данных: user_7d_purchase_count, item_avg_rating, session_clicks_last_hour. Модель обучается на одних данных (offline, батчем по историческим логам) и делает предсказания на других (online, real-time по текущему состоянию).
Если формулы расчёта фичей для обучения и инференса расходятся -- получаем training/serving skew. Типичные причины:
- Дубль логики: SQL-запрос для обучения и Python-функция для онлайна -- пишутся разными людьми, расходятся
- Разный timestamp: обучение считает
user_7d_countна момент события, продакшн -- на момент запроса - Разные источники: батч читает из data lake, онлайн -- из PostgreSQL, там другая схема
- Backfill-несовместимость: новую фичу нельзя пересчитать на старых данных, появляется NaN
Результат: в offline-метриках модель хороша, в проде -- деградирует. Отлаживать такое больно.
Feature store решает эту проблему как централизованный источник фичей с единым определением, из которого получают данные и offline-pipeline, и online-serving.
Архитектура feature store
┌──────────────────────────────────────────────┐
│ Feature Definitions (code) │
│ - entity definitions │
│ - transformations (SQL / UDF / streaming) │
│ - schemas │
└───────────┬───────────────────────┬──────────┘
│ materialize │ sync
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ Offline store │──────>│ Online store │
│ Parquet/Hive/ │ │ Redis/DynamoDB │
│ BigQuery/S3 │ │ Cassandra │
│ -- historical │ │ -- latest │
│ -- point-in-time│ │ -- low latency │
└─────────┬────────┘ └────────┬─────────┘
│ │
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ Training pipeline│ │ Serving API │
│ (Spark / Flink) │ │ (<10 ms) │
└──────────────────┘ └──────────────────┘
Offline store
Для обучения и backfill. Хранит историю фичей с timestamp'ами. Формат: Parquet / Delta Lake / Iceberg в S3 / HDFS, или таблицы в BigQuery / Snowflake / Hive.
Запросы: "дай значения фичей для этого набора (user_id, timestamp) пар -- ровно те значения, что были на тот момент". Это point-in-time join.
Online store
Для inference в реальном времени. Хранит только последние значения фичей.
Требования: латентность < 10 мс p99, поддержка batch-get (запрос фичей для сотен item-ов за раз для ранжирования). Типичные выборы: Redis (простота, низкая латентность), DynamoDB (managed, автоскейл), Cassandra (масштаб), Aerospike (сверхнизкая латентность).
Ключ: обычно entity_id + feature_view_name → значение или map.
Материализация
Процесс "упаковки" фичей из offline-store в online-store. Запускается по расписанию (раз в N минут/час) или по событию.
Offline store (Parquet)
│
│ materialize job (Spark)
▼
Online store (Redis)
user:42:v_last_7d -> {purchases: 3, revenue: 1240.50}
Для streaming features (обновляются событиями) вместо материализации используется прямой stream в online store через Kafka → Flink → Redis.
Point-in-time correctness
Это главная техническая фича feature store. Представьте набор меток:
(user=42, ts=2025-08-15 10:00, label=1)
(user=42, ts=2025-09-02 14:00, label=0)
(user=17, ts=2025-08-20 09:00, label=1)
Нужна фича user_7d_purchases. Для каждой метки -- значение фичи на момент timestamp, не последнее. Если для user=42 на 2025-08-15 было 5 покупок за 7 дней, а на 2025-09-02 -- 2, нужны именно эти значения, а не последнее "3".
Наивный LEFT JOIN по user_id даст data leakage -- модель увидит будущее. Feature store делает правильный point-in-time as-of join:
-- Simplified point-in-time join for one feature
SELECT
l.user_id,
l.event_ts,
l.label,
f.purchases_7d
FROM labels l
LEFT JOIN LATERAL (
SELECT purchases_7d
FROM feature_user_activity fu
WHERE fu.user_id = l.user_id
AND fu.feature_ts <= l.event_ts -- "as-of": only past
AND fu.feature_ts > l.event_ts - INTERVAL '1 day' -- TTL guard
ORDER BY fu.feature_ts DESC
LIMIT 1
) f ON TRUE;
Feature store оборачивает такой join в декларативный API -- пользователь просто говорит "хочу эти фичи для этих (entity, ts) пар", а движок строит корректный SQL / Spark-job.
Feature definitions as code
Вместо SQL-в-ноутбуке фичи определяются как типизированные декларации. Пример в стиле Feast (инструмент №1 в open source):
from datetime import timedelta
from feast import Entity, FeatureView, FeatureService, Field
from feast.types import Float32, Int64
from feast.infra.offline_stores.file_source import FileSource
# 1. Entity -- business key
user = Entity(
name="user",
join_keys=["user_id"],
description="Registered marketplace user",
)
# 2. Source -- where raw data lives
user_activity_source = FileSource(
name="user_activity",
path="s3://dwh/feature_source/user_activity/", # partitioned Parquet
timestamp_field="event_timestamp",
created_timestamp_column="ingestion_timestamp",
)
# 3. Feature View -- logical group of features
user_activity_fv = FeatureView(
name="user_activity_7d",
entities=[user],
ttl=timedelta(days=3), # stale threshold for online serving
schema=[
Field(name="purchases_7d", dtype=Int64),
Field(name="revenue_7d", dtype=Float32),
Field(name="sessions_7d", dtype=Int64),
Field(name="avg_session_sec", dtype=Float32),
],
source=user_activity_source,
online=True, # materialize to online store
)
# 4. Feature Service -- named bundle for a model
recommender_v3 = FeatureService(
name="recommender_v3",
features=[
user_activity_fv,
user_activity_fv[["purchases_7d"]], # subset allowed
],
)
Благодаря декларативному определению:
- Тот же объект
user_activity_fvчитают и offline-training, и online-serving - Определения версионируются в git
- CI проверяет схему, запускает dry-run
- Lineage автоматический: от label до источника данных
Материализация и пайплайны
Типичный цикл:
- Raw events пишутся в Kafka или batch в S3 (
/events/...) - Transformation job (Spark/Flink/dbt) агрегирует в фичи:
user_7d_purchases = COUNT(*) WHERE ts > now()-7d - Результат пишется в offline store как Parquet-партиция по дате
- Materialization job (feature store) раз в N минут копирует последние значения из offline → online store
- Online serving API отдаёт фичи за < 10 мс
Для streaming-фичей (обновляются событиями, например last_item_viewed) -- отдельная ветка: Kafka → Flink → online store напрямую, без offline. Но offline-копия всё равно пишется для обучения (можно через CDC из online).
Tooling сравнение
| Инструмент | Тип | Offline | Online | Streaming | Pricing |
|---|---|---|---|---|---|
| Feast | OSS | любой SQL / Parquet | Redis, DynamoDB, BigTable | ограничено | free |
| Tecton | Managed | Spark / Snowflake | DynamoDB, Redis | first-class | enterprise |
| Hopsworks | OSS + managed | own feature store | RonDB | Spark streaming | free / commercial |
| Databricks FS | Managed | Delta Lake | Azure/AWS backends | Structured Streaming | part of Databricks |
| Vertex AI FS | Managed (GCP) | BigQuery | BigTable | streaming API | pay-per-use |
| SageMaker FS | Managed (AWS) | S3 | DynamoDB-подобное | Kinesis | pay-per-use |
| Chronon (LinkedIn) | OSS | Hive | KV store | Flink | free |
Когда Feast: хотите OSS, уже есть data lake, streaming не в приоритете. Лёгкая интеграция со Spark, Python-first API.
Когда Tecton: готовы платить за managed, нужен streaming как first-class, сложные on-demand transformations, требуется governance.
Когда cloud-native (Vertex/SageMaker): уже в экосистеме, не хотите содержать инфру.
Пример фичей для рекомендательной системы
Фичи для рекомендательного движка маркетплейса:
User features (entity=user_id, TTL=1 day):
purchases_7d,purchases_30d,purchases_lifetimeavg_order_value_30dcategories_viewed_1d(list of category_ids)sessions_1d,session_avg_duration_secdevice_type,preferred_language
Item features (entity=item_id, TTL=6 hours):
views_24h,purchases_24h,add_to_cart_24hctr_24h(click-through rate)price,discount_pct,stock_qtycategory_id,brand_id,tags(list)avg_rating,review_count
Context features (no entity, computed on-demand):
hour_of_day,day_of_weekdevice_type(из request)session_item_views(streaming, из текущей сессии)
On-demand features (transformations at request time):
user_item_category_affinity=user_category_views[item.category_id] / user_total_views
При ранжировании модель получает для каждой пары (user, item_candidate) объединённый вектор из user + item + context + on-demand фичей. В Feast это один вызов:
features = store.get_online_features(
features=[
"user_activity_7d:purchases_7d",
"user_activity_7d:revenue_7d",
"item_stats_24h:ctr_24h",
"item_stats_24h:views_24h",
],
entity_rows=[
{"user_id": 42, "item_id": 1001},
{"user_id": 42, "item_id": 1002},
# ... 200 candidates
],
).to_dict()
Feature retrieval из приложения
Python (Feast)
# Online serving integration in a FastAPI app
from feast import FeatureStore
from fastapi import FastAPI
import time
import structlog
app = FastAPI()
store = FeatureStore(repo_path="./feature_repo")
log = structlog.get_logger()
@app.post("/rank")
async def rank_items(user_id: int, candidates: list[int]):
started = time.monotonic()
# Batch fetch features for all (user, item) pairs
entity_rows = [{"user_id": user_id, "item_id": c} for c in candidates]
features = store.get_online_features(
features=[
"user_activity_7d:purchases_7d",
"user_activity_7d:revenue_7d",
"item_stats_24h:ctr_24h",
"item_stats_24h:price",
"item_stats_24h:avg_rating",
],
entity_rows=entity_rows,
).to_dict()
# Pass to model, get scores
scores = model.predict(features)
ranked = sorted(zip(candidates, scores), key=lambda x: -x[1])
log.info(
"rank.completed",
user_id=user_id,
candidates=len(candidates),
latency_ms=int((time.monotonic() - started) * 1000),
)
return [item for item, _ in ranked[:20]]
Go (custom online store)
Не все команды используют Feast -- часто feature store делают in-house поверх Redis. Ключевое -- сохранить контракт: batch-get по списку entities.
package features
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"time"
"github.com/redis/go-redis/v9"
)
// FeatureView defines a named group of features for an entity type.
type FeatureView struct {
Name string
EntityID string // "user_id", "item_id"
TTL time.Duration // staleness bound
Fields []string // feature names
}
// UserActivity7d is the Go-side representation of the feature view.
type UserActivity7d struct {
Purchases7d int64 `json:"purchases_7d"`
Revenue7d float64 `json:"revenue_7d"`
Sessions7d int64 `json:"sessions_7d"`
AvgSessionSec float32 `json:"avg_session_sec"`
}
// OnlineStore fetches features from Redis with batched MGET.
type OnlineStore struct {
rdb *redis.Client
log *slog.Logger
}
func NewOnlineStore(rdb *redis.Client, log *slog.Logger) *OnlineStore {
return &OnlineStore{rdb: rdb, log: log}
}
// GetUserActivity fetches features for multiple users in one round-trip.
// Returns a map indexed by user_id. Missing users get zero-value with warning.
func (s *OnlineStore) GetUserActivity(
ctx context.Context,
userIDs []int64,
) (map[int64]UserActivity7d, error) {
if len(userIDs) == 0 {
return nil, nil
}
keys := make([]string, len(userIDs))
for i, id := range userIDs {
keys[i] = fmt.Sprintf("fv:user_activity_7d:%d", id)
}
started := time.Now()
raw, err := s.rdb.MGet(ctx, keys...).Result()
if err != nil {
return nil, fmt.Errorf("mget user activity: %w", err)
}
out := make(map[int64]UserActivity7d, len(userIDs))
missing := 0
for i, v := range raw {
if v == nil {
missing++
out[userIDs[i]] = UserActivity7d{} // cold-start defaults
continue
}
var fv UserActivity7d
if err := json.Unmarshal([]byte(v.(string)), &fv); err != nil {
s.log.Warn("feature.unmarshal_failed", "user_id", userIDs[i], "err", err)
continue
}
out[userIDs[i]] = fv
}
s.log.Info("feature.fetch",
"view", "user_activity_7d",
"requested", len(userIDs),
"missing", missing,
"latency_ms", time.Since(started).Milliseconds(),
)
return out, nil
}
Governance и версионирование
Фичи -- продукт. У них должны быть:
- Владелец (team / on-call)
- SLA на freshness и availability
- Lineage -- откуда берутся (какие источники, какие таблицы upstream)
- Документация -- описание, единицы измерения, ожидаемый диапазон
- Тесты -- проверка schema, выбросов, null-rate
Feast / Tecton генерируют lineage-графы автоматически. Deprecated-фичи помечаются, но не удаляются сразу -- модели могут их использовать.
Версионирование: обычно два подхода --
- Неизменяемое определение: новая версия = новое имя (
user_activity_v2) - Versioned features (Tecton):
user_activity@v3, совместимость в пределах major-версии
Observability
Метрики, без которых feature store в проде не живёт:
feature_store_online_latency_msp50/p99 по feature_viewfeature_store_staleness_seconds(current_time - feature_ts) -- критичноfeature_store_null_rate-- сколько сущностей без фичи (cold start)feature_store_materialization_duration-- длительность job'овfeature_store_training_serving_skew-- периодический recon: сравнить несколько выборок онлайн vs офлайн значения
Алерты: staleness > TTL, null_rate > baseline + 3σ, materialization failures.
Когда feature store избыточен
Feature store -- серьёзная инфраструктура. Если у вас:
- 1-2 модели, фичи посчитаны по одной таблице в том же Postgres
- Нет streaming features
- Offline и online пишет одна команда
-- feature store только добавит операционные издержки. Достаточно материализованных view в Postgres и явных версионированных ETL-джоб. Порог внедрения обычно -- когда моделей 5+, или когда появляется streaming, или когда несколько команд делят фичи.
Выводы
- Training/serving skew -- главная причина деградации ML в проде. Feature store -- центральное решение этой проблемы.
- Архитектура: feature definitions (как код) → offline store (исторические, Parquet) + online store (последние значения, Redis/DynamoDB).
- Point-in-time correctness через as-of joins -- ключевая техническая возможность. Без неё -- data leakage.
- Feast -- стандарт OSS, Tecton -- managed с фокусом на streaming, cloud-native (Vertex/SageMaker) -- если уже в экосистеме.
- Feature retrieval -- batched по сотням entities за один round-trip, p99 < 10 мс.
- Инструмент не бесплатен по сложности: вводите, когда есть множество моделей, streaming или кросс-командный шаринг фичей.