从分析到产品

在前面的文章中,我们做了大量分析工作:清洗数据、可视化趋势、建立模型、挖掘洞察。但这些分析的产出大多是 Jupyter Notebook 或 PDF 报告,业务团队无法直接交互使用。

数据产品化就是把分析成果转化为可交互、可复用、可部署的工具。典型的数据产品包括:

  • 交互式仪表盘:业务人员可以自助查看指标
  • API 服务:其他系统可以调用模型预测
  • 自动化报告:定时生成并分发
  • 数据目录:帮助团队理解和发现数据

用 Streamlit 构建交互式仪表盘

Streamlit 是 Python 生态中最流行的数据应用框架。它让你用纯 Python 写出漂亮的 Web 应用,不需要 HTML/CSS/JS。

安装

pip install streamlit pandas matplotlib plotly

基础仪表盘示例

# dashboard.py
import streamlit as st
import pandas as pd
import plotly.express as px
import plotly.graph_objects as go
from datetime import datetime, timedelta

# 页面配置
st.set_page_config(
    page_title="电商销售仪表盘",
    page_icon="📊",
    layout="wide"
)

st.title("📊 电商销售实时仪表盘")
st.markdown("---")

# 缓存数据加载(避免每次交互都重读)
@st.cache_data
def load_data():
    df = pd.read_csv('sales_data.csv', parse_dates=['order_date'])
    return df

df = load_data()

# 侧边栏过滤
st.sidebar.header("过滤条件")

date_range = st.sidebar.date_input(
    "日期范围",
    value=(df['order_date'].min(), df['order_date'].max()),
    min_value=df['order_date'].min(),
    max_value=df['order_date'].max()
)

selected_categories = st.sidebar.multiselect(
    "产品品类",
    options=df['category'].unique(),
    default=df['category'].unique()[:3]
)

selected_city = st.sidebar.selectbox(
    "城市",
    options=['全部'] + sorted(df['city'].unique().tolist())
)

# 日期过滤
mask = (df['order_date'] >= pd.to_datetime(date_range[0])) & \
       (df['order_date'] <= pd.to_datetime(date_range[1])) & \
       (df['category'].isin(selected_categories))

if selected_city != '全部':
    mask &= (df['city'] == selected_city)

filtered_df = df[mask]

# 核心指标
col1, col2, col3, col4 = st.columns(4)

with col1:
    total_revenue = filtered_df['amount'].sum()
    st.metric(
        label="总销售额",
        value=f{total_revenue:,.0f}",
        delta="5.2%"  # 可对接真实同比
    )

with col2:
    total_orders = len(filtered_df)
    st.metric(label="总订单数", value=f"{total_orders:,}")

with col3:
    avg_order = filtered_df['amount'].mean()
    st.metric(label="平均客单价", value=f{avg_order:.2f}")

with col4:
    unique_customers = filtered_df['customer_id'].nunique()
    st.metric(label="活跃客户数", value=f"{unique_customers:,}")

st.markdown("---")

# 双栏布局
left_col, right_col = st.columns(2)

with left_col:
    st.subheader("每日销售趋势")
    
    daily_sales = filtered_df.groupby(
        filtered_df['order_date'].dt.date
    )['amount'].sum().reset_index()
    
    fig = px.line(
        daily_sales, x='order_date', y='amount',
        markers=True, line_shape='spline'
    )
    fig.update_layout(
        xaxis_title="日期",
        yaxis_title="销售额",
        hovermode='x unified'
    )
    st.plotly_chart(fig, use_container_width=True)

with right_col:
    st.subheader("品类销售占比")
    
    category_sales = filtered_df.groupby('category')['amount'].sum().reset_index()
    fig = px.pie(category_sales, values='amount', names='category')
    st.plotly_chart(fig, use_container_width=True)

# 第二行
st.subheader("各城市销售排行")
city_sales = filtered_df.groupby('city')['amount'].sum() \
    .sort_values(ascending=True).reset_index()

fig = px.bar(
    city_sales.tail(10), x='amount', y='city',
    orientation='h', color='amount',
    color_continuous_scale='Blues'
)
fig.update_layout(xaxis_title="销售额", yaxis_title="城市")
st.plotly_chart(fig, use_container_width=True)

# 数据表格(可折叠)
with st.expander("查看原始数据"):
    st.dataframe(filtered_df.head(1000), use_container_width=True)

运行仪表盘

streamlit run dashboard.py

Streamlit 的常用组件

import streamlit as st

# 1. 文本
st.title("标题")
st.header("章节标题")
st.subheader("子标题")
st.markdown("支持 **Markdown** 格式")
st.caption("小字说明")
st.latex(r"E = mc^2")
st.code("print('hello')", language='python')

# 2. 数据展示
st.dataframe(df)  # 交互式表格
st.table(df.head())  # 静态表格
st.json({"key": "value"})

# 3. 图表
st.line_chart(data)
st.bar_chart(data)
st.area_chart(data)
st.map(geo_df)

# 4. 输入控件
st.button("点击")
st.download_button("下载", data)
st.checkbox("勾选")
st.radio("单选", options=['A', 'B', 'C'])
st.selectbox("下拉选择", options=['X', 'Y', 'Z'])
st.multiselect("多选", options=[1, 2, 3])
st.slider("滑块", 0, 100, 50)
st.text_input("文本输入")
st.text_area("多行文本")
st.date_input("日期选择")
st.file_uploader("文件上传")

# 5. 布局
st.sidebar.markdown("侧边栏")
col1, col2, col3 = st.columns(3)
with col1:
    st.write("第一列")
tab1, tab2 = st.tabs(["标签1", "标签2"])
with tab1:
    st.write("第一个标签页")

# 6. 状态
st.success("成功")
st.info("信息")
st.warning("警告")
st.error("错误")
st.spinner("加载中...")
st.progress(0.5)
st.balloons()  # 🎈

# 7. 缓存
@st.cache_data
def expensive_computation():
    # 耗时操作,结果会被缓存
    return result

@st.cache_resource
def load_model():
    # 缓存模型等全局资源
    return model

部署 Streamlit 应用

# 本地部署
streamlit run app.py --server.port 8501

# Docker 部署
docker run -p 8501:8501 \
  -v $(pwd)/app:/app \
  streamlit/app:latest

# Streamlit Cloud(免费)
# 1. 把代码推到 GitHub
# 2. 在 share.streamlit.io 上连接仓库
# 3. 自动部署

用 FastAPI 构建模型 API

把机器学习模型包装成 REST API,让其他系统调用。

安装

pip install fastapi uvicorn pydantic scikit-learn

模型服务示例

# model_api.py
import pickle
import numpy as np
import pandas as pd
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field
from typing import Optional
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

app = FastAPI(
    title="客户价值预测 API",
    description="预测客户价值的 REST 服务",
    version="1.0.0"
)

# 定义请求 Schema
class CustomerFeatures(BaseModel):
    recency: int = Field(..., ge=0, description="最近购买天数")
    frequency: int = Field(..., ge=1, description="购买次数")
    monetary: float = Field(..., ge=0, description="总消费金额")
    avg_order_value: float = Field(..., ge=0, description="平均客单价")
    days_since_registration: int = Field(..., ge=0, description="注册天数")
    category_diversity: int = Field(..., ge=1, description="购买品类数")

class PredictionResponse(BaseModel):
    customer_id: Optional[str] = None
    predicted_ltv: float
    customer_segment: str
    confidence: float

# 加载模型(启动时一次加载)
try:
    with open('models/ltv_model.pkl', 'rb') as f:
        model = pickle.load(f)
    with open('models/scaler.pkl', 'rb') as f:
        scaler = pickle.load(f)
    logger.info("模型加载成功")
except FileNotFoundError:
    model = None
    scaler = None
    logger.warning("模型文件未找到,使用模拟模式")

@app.get("/")
def root():
    return {
        "service": "客户价值预测 API",
        "status": "running",
        "model_loaded": model is not None
    }

@app.get("/health")
def health_check():
    return {"status": "healthy"}

@app.post("/predict", response_model=PredictionResponse)
def predict(features: CustomerFeatures):
    """预测客户价值"""
    
    # 输入验证
    if features.monetary < 0:
        raise HTTPException(status_code=400, detail="消费金额不能为负")
    
    if model is None:
        # 模拟预测
        predicted_ltv = features.monetary * 1.5 + features.frequency * 100
        segment = "高价值" if predicted_ltv > 5000 else "普通"
        confidence = 0.75
    else:
        # 真实预测
        feature_array = np.array([[
            features.recency,
            features.frequency,
            features.monetary,
            features.avg_order_value,
            features.days_since_registration,
            features.category_diversity
        ]])
        
        if scaler:
            feature_array = scaler.transform(feature_array)
        
        predicted_ltv = model.predict(feature_array)[0]
        # 计算置信度(简化示例)
        confidence = min(0.95, max(0.5, predicted_ltv / 10000))
        
        # 分群
        if predicted_ltv > 8000:
            segment = "VIP"
        elif predicted_ltv > 3000:
            segment = "高价值"
        elif predicted_ltv > 1000:
            segment = "中等价值"
        else:
            segment = "低价值"
    
    logger.info(
        f"预测: LTV={predicted_ltv:.2f}, 分群={segment}, 置信度={confidence:.2f}"
    )
    
    return PredictionResponse(
        predicted_ltv=round(predicted_ltv, 2),
        customer_segment=segment,
        confidence=round(confidence, 3)
    )

@app.post("/batch_predict")
def batch_predict(features_list: list[CustomerFeatures]):
    """批量预测"""
    results = []
    for features in features_list:
        result = predict(features)
        results.append(result)
    return {"predictions": results, "count": len(results)}

# 启动命令
# uvicorn model_api:app --reload --host 0.0.0.0 --port 8000

测试 API

# 启动服务
uvicorn model_api:app --reload --port 8000

# 请求预测
curl -X POST "http://localhost:8000/predict" \
  -H "Content-Type: application/json" \
  -d '{
    "recency": 5,
    "frequency": 12,
    "monetary": 5000.0,
    "avg_order_value": 416.67,
    "days_since_registration": 365,
    "category_diversity": 4
  }'

# 批量预测
curl -X POST "http://localhost:8000/batch_predict" \
  -H "Content-Type: application/json" \
  -d '[
    {"recency": 5, "frequency": 12, "monetary": 5000, "avg_order_value": 416.67, "days_since_registration": 365, "category_diversity": 4},
    {"recency": 30, "frequency": 3, "monetary": 800, "avg_order_value": 266.67, "days_since_registration": 180, "category_diversity": 2}
  ]'

自动生成 API 文档

FastAPI 自动生成 OpenAPI 文档,访问以下地址查看:

  • Swagger UI: http://localhost:8000/docs
  • ReDoc: http://localhost:8000/redoc

用 Papermill 自动化报告

Papermill 可以参数化执行 Jupyter Notebook,生成定制化报告。

安装

pip install papermill

创建参数化 Notebook

在 notebook 的第一个 cell 中定义参数:

# 参数(会被 Papermill 覆写)
report_date = "2025-01-23"
segment = "全部"
region = "华东"

命令行运行

# 用参数执行 notebook
papermill \
  report_template.ipynb \
  output/report_2025-01-23.ipynb \
  -p report_date "2025-01-23" \
  -p segment "VIP" \
  -p region "华东"

# 自动导出为 HTML
jupyter nbconvert \
  --to html \
  output/report_2025-01-23.ipynb \
  --output report_2025-01-23.html

批量生成报告

import papermill as pm
from datetime import datetime, timedelta

# 批量生成最近7天的报告
today = datetime.now()

for i in range(1, 8):
    report_date = (today - timedelta(days=i)).strftime('%Y-%m-%d')
    
    for segment in ['全部', 'VIP', '普通']:
        input_path = 'notebooks/daily_report_template.ipynb'
        output_path = f'reports/{report_date}_{segment}_report.ipynb'
        
        try:
            pm.execute_notebook(
                input_path,
                output_path,
                parameters={
                    'report_date': report_date,
                    'segment': segment,
                    'generate_charts': True,
                    'send_email': (segment == '全部')  # 仅对全量报告发邮件
                }
            )
            print(f'✅ 报告生成成功: {output_path}')
        except Exception as e:
            print(f'❌ 报告生成失败: {output_path}, 错误: {e}')

数据文档与治理

数据字典

数据字典是数据产品的基础。它让所有人都能理解每个字段的含义。

# data_dictionary.py
DATA_DICTIONARY = {
    "orders": {
        "description": "订单主表",
        "update_frequency": "实时",
        "owner": "data_team@company.com",
        "fields": {
            "order_id": {
                "type": "string",
                "description": "订单唯一标识",
                "example": "ORD-20250101-00001",
                "constraints": "PRIMARY KEY, NOT NULL"
            },
            "customer_id": {
                "type": "string",
                "description": "客户编号",
                "example": "CUST-12345",
                "constraints": "FOREIGN KEY → customers.customer_id"
            },
            "amount": {
                "type": "float",
                "description": "订单金额(元)",
                "example": 299.00,
                "constraints": "> 0"
            },
            "status": {
                "type": "string",
                "description": "订单状态",
                "enum": ["pending", "completed", "cancelled", "returned"],
                "example": "completed"
            },
            "order_date": {
                "type": "datetime",
                "description": "下单时间",
                "example": "2025-01-23 14:30:00"
            }
        }
    },
    "customers": {
        "description": "客户信息表",
        "update_frequency": "每日更新",
        "owner": "customer_team@company.com",
        "fields": {
            "customer_id": {
                "type": "string",
                "description": "客户唯一标识",
                "example": "CUST-12345",
                "constraints": "PRIMARY KEY"
            },
            "registration_date": {
                "type": "date",
                "description": "注册日期",
                "example": "2024-06-15"
            },
            "city": {
                "type": "string",
                "description": "所在城市",
                "example": "北京"
            },
            "segment": {
                "type": "string",
                "description": "客户分群",
                "enum": ["VIP", "高价值", "中等", "低价值"],
                "example": "VIP"
            }
        }
    }
}

DVC(数据版本控制)

像管理代码一样管理数据。

# 安装
pip install dvc

# 初始化
dvc init

# 添加远程存储(比如 S3、GCS、阿里云OSS)
dvc remote add -d myremote s3://my-data-bucket/dvc-store

# 追踪数据文件
dvc add data/sales_2025.parquet

# 提交到 git
git add data/sales_2025.parquet.dvc .gitignore
git commit -m "feat: track sales data with DVC"
git push

# 推送数据到远程存储
dvc push

# 在另一台机器上拉取数据
dvc pull

MLOps 基础

MLOps 是把机器学习模型从实验阶段带到生产环境的一套实践。

核心组件

# model_registry.py - 简单的模型注册表
import mlflow
import mlflow.sklearn
from datetime import datetime

class ModelRegistry:
    """模型注册表"""
    
    def __init__(self, tracking_uri='sqlite:///mlflow.db'):
        mlflow.set_tracking_uri(tracking_uri)
    
    def log_model(self, model, params, metrics, model_name):
        """记录模型实验"""
        with mlflow.start_run(run_name=f"{model_name}_{datetime.now():%Y%m%d}"):
            # 记录参数
            for key, value in params.items():
                mlflow.log_param(key, value)
            
            # 记录指标
            for key, value in metrics.items():
                mlflow.log_metric(key, value)
            
            # 记录模型
            mlflow.sklearn.log_model(model, model_name)
            
            return mlflow.active_run().info.run_id
    
    def load_production_model(self, model_name):
        """加载生产环境模型"""
        model_uri = f"models:/{model_name}/Production"
        return mlflow.sklearn.load_model(model_uri)
    
    def promote_to_production(self, model_name, version):
        """将模型提升到生产环境"""
        client = mlflow.tracking.MlflowClient()
        client.transition_model_version_stage(
            name=model_name,
            version=version,
            stage="Production"
        )

特征存储(Feature Store)

# feature_store.py - 简单的特征存储
import redis
import json

class FeatureStore:
    """在线特征存储"""
    
    def __init__(self, host='localhost', port=6379):
        self.client = redis.Redis(host=host, port=port, decode_responses=True)
    
    def set_features(self, customer_id, features, ttl=86400):
        """存储客户特征"""
        key = f"features:customer:{customer_id}"
        self.client.setex(key, ttl, json.dumps(features))
    
    def get_features(self, customer_id):
        """获取客户特征"""
        key = f"features:customer:{customer_id}"
        data = self.client.get(key)
        return json.loads(data) if data else None
    
    def batch_get_features(self, customer_ids):
        """批量获取"""
        pipeline = self.client.pipeline()
        for cid in customer_ids:
            pipeline.get(f"features:customer:{cid}")
        results = pipeline.execute()
        return [json.loads(r) if r else None for r in results]

监控数据漂移

# monitor_drift.py - 数据漂移监控
from scipy.stats import ks_2samp
import numpy as np

def detect_drift(reference_data, new_data, threshold=0.05):
    """检测数据分布漂移"""
    drift_report = {}
    
    for column in reference_data.columns:
        if reference_data[column].dtype in ['int64', 'float64']:
            # 数值型:KS 检验
            ks_stat, p_value = ks_2samp(
                reference_data[column].dropna(),
                new_data[column].dropna()
            )
            drift_report[column] = {
                'test': 'KS',
                'statistic': ks_stat,
                'p_value': p_value,
                'drift_detected': p_value < threshold
            }
    
    drift_count = sum(1 for v in drift_report.values() if v['drift_detected'])
    total = len(drift_report)
    
    print(f"漂移检测完成: {drift_count}/{total} 个特征检测到漂移")
    
    if drift_count > total * 0.3:
        print("⚠️ 严重警告:超过 30% 的特征发生漂移,建议重新训练模型")
    
    return drift_report

用数据讲故事

数据产品化不只是技术实现,还包括如何有效沟通分析结果。

好的数据故事结构

  1. 设定背景:我们在关注什么问题?为什么这个问题重要?
  2. 展示洞察:数据告诉我们什么?用图表和数字说话
  3. 解释原因:为什么会出现这种情况?可能的根因是什么?
  4. 提出建议:基于洞察,我们建议做什么?
  5. 量化影响:如果采纳建议,预期能带来什么效果?

数据展示的三个原则

1. 先结论后细节

不要像侦探小说一样把结论放在最后。先说结论,再展开论证。

# ❌ 错误的展示顺序
# "过去三个月中,我们的销售额从 100 万增长到 120 万,
#  再到 150 万,期间发生了 X、Y、Z 事件,综合来看
#  是因为我们的营销策略调整有效,因此建议继续加大投入。"

# ✅ 正确的展示顺序
# "建议继续加大营销投入。过去三个月销售额持续增长20%。
#  数据表明增长主要来自新渠道投放,ROI 达到 3.5。
#  具体增长曲线和事件分析如下..."

2. 一个图表只讲一个故事

一个图表有多个信息时,读者的注意力会被分散。

3. 选择合适的图表类型

意图 推荐图表
趋势变化 折线图
对比排序 柱状图
分布情况 直方图、箱线图
组成比例 饼图、堆叠柱状图
相关性 散点图
地理分布 地图

小结

数据产品化让数据分析的价值从个人扩展到整个组织。无论是指标仪表盘、模型 API,还是自动化报告,它们都让数据更易获取、更可行动。

本篇文章的核心知识点:

  • 用 Streamlit 快速构建交互式数据应用
  • 用 FastAPI 将模型包装为 RESTful API
  • 用 Papermill 自动化生成定制化报告
  • 数据字典和 DVC 实现数据治理
  • MLOps 的模型注册、特征存储和漂移监控
  • 数据故事讲述的结构和原则

在下一篇文章中,也是本系列的最后一篇,我们将回顾整个学习旅程,并为你的持续成长规划一条清晰的路径。