ML feature platform

From product DBs to online and offline features

Dataflow13 nodes · 13 connectionsWritten by an agent

diagram.html#The link keeps the view, selection, route and playback.

What the agent was asked

“Show how data flows through our ML feature platform. Product databases (Postgres) are captured with Debezium CDC into Kafka. A Flink job computes streaming features and writes them to Redis (online store) and to a Delta Lake on S3 (offline store). Airflow runs nightly Spark jobs over the lake to build training sets, which a training pipeline on SageMaker consumes; models are registered in MLflow. The prediction service reads online features from Redis and loads models from MLflow. Also note that raw Kafka topics are archived to S3 for replay.”

Written by a coding agent in the skill’s eval: it read the request, wrote the JSON, validated it and delivered the file.

Read about dataflow diagrams

The JSON

ml-feature-platform/diagram.json233 lines
{  "kind": "dataflow",  "title": "ML feature platform",  "subtitle": "From product DBs to online and offline features",  "direction": "DOWN",  "groups": [    { "id": "source", "label": "Source" },    { "id": "cdc", "label": "CDC & streaming" },    { "id": "stores", "label": "Feature stores" },    { "id": "archive", "label": "Archive" },    { "id": "batch", "label": "Batch training" },    { "id": "training", "label": "Training & registry" },    { "id": "serving", "label": "Serving" }  ],  "nodes": [    {      "id": "postgres",      "type": "database",      "group": "source",      "card": {        "title": "Product DBs",        "subtitle": "PostgreSQL",        "brand": "postgresql"      }    },    {      "id": "debezium",      "type": "service",      "group": "cdc",      "card": { "title": "Debezium", "subtitle": "CDC connector" }    },    {      "id": "kafka",      "type": "queue",      "group": "cdc",      "card": { "title": "Kafka", "subtitle": "Raw change topics" }    },    {      "id": "flink",      "type": "service",      "group": "cdc",      "card": { "title": "Flink job", "subtitle": "Streaming features" }    },    {      "id": "redis",      "type": "cache",      "group": "stores",      "card": {        "title": "Redis",        "subtitle": "Online feature store",        "brand": "redis"      }    },    {      "id": "delta-lake",      "type": "storage",      "group": "stores",      "card": {        "title": "Delta Lake",        "subtitle": "Offline store on S3"      }    },    {      "id": "kafka-archive",      "type": "storage",      "group": "archive",      "card": {        "title": "Topic archive",        "subtitle": "Raw Kafka on S3",        "footer": { "left": { "text": "Replay" } }      }    },    {      "id": "airflow",      "type": "service",      "group": "batch",      "card": {        "title": "Airflow",        "subtitle": "Nightly orchestration"      }    },    {      "id": "spark",      "type": "service",      "group": "batch",      "card": {        "title": "Spark job",        "subtitle": "Builds training sets"      }    },    {      "id": "training-sets",      "type": "storage",      "group": "batch",      "card": { "title": "Training sets", "subtitle": "Curated on S3" }    },    {      "id": "sagemaker",      "type": "service",      "group": "training",      "card": { "title": "SageMaker", "subtitle": "Training pipeline" }    },    {      "id": "mlflow",      "type": "service",      "group": "training",      "card": { "title": "MLflow", "subtitle": "Model registry" }    },    {      "id": "prediction-service",      "type": "service",      "group": "serving",      "card": {        "title": "Prediction service",        "subtitle": "Online inference"      }    }  ],  "edges": [    {      "id": "postgres-debezium",      "from": "postgres",      "to": "debezium",      "kind": "async",      "label": "WAL stream"    },    {      "id": "debezium-kafka",      "from": "debezium",      "to": "kafka",      "kind": "async",      "label": "publish"    },    {      "id": "kafka-flink",      "from": "kafka",      "to": "flink",      "kind": "async",      "label": "consume"    },    {      "id": "kafka-archive-edge",      "from": "kafka",      "to": "kafka-archive",      "kind": "async",      "label": "archive raw"    },    {      "id": "flink-redis",      "from": "flink",      "to": "redis",      "label": "write features"    },    {      "id": "flink-delta",      "from": "flink",      "to": "delta-lake",      "label": "write features"    },    {      "id": "delta-spark",      "from": "delta-lake",      "to": "spark",      "label": "read nightly"    },    {      "id": "airflow-spark",      "from": "airflow",      "to": "spark",      "label": "trigger"    },    {      "id": "spark-trainingsets",      "from": "spark",      "to": "training-sets",      "label": "build sets"    },    {      "id": "trainingsets-sagemaker",      "from": "training-sets",      "to": "sagemaker",      "label": "train"    },    {      "id": "sagemaker-mlflow",      "from": "sagemaker",      "to": "mlflow",      "label": "register"    },    {      "id": "redis-prediction",      "from": "redis",      "to": "prediction-service",      "label": "read features"    },    {      "id": "mlflow-prediction",      "from": "mlflow",      "to": "prediction-service",      "label": "load model"    }  ],  "views": [    {      "id": "online-path",      "label": "Online path",      "caption": "CDC to online serving via Redis",      "nodes": [        "postgres",        "debezium",        "kafka",        "flink",        "redis",        "prediction-service"      ]    },    {      "id": "offline-path",      "label": "Offline path",      "caption": "Lake to training to model registry",      "nodes": [        "flink",        "delta-lake",        "airflow",        "spark",        "training-sets",        "sagemaker",        "mlflow",        "prediction-service"      ]    }  ]}

More dataflow examples

All examples