Skip to content

A12 前端数据工程

目标:掌握前端数据埋点、AB 实验、指标体系和数据看板架构,用数据驱动产品和技术决策。


核心要点(TL;DR)

  • 前端数据工程是连接用户行为与业务决策的桥梁。
  • 埋点体系需要标准化、可治理、可验证。
  • AB 实验是验证产品假设的科学方法,需要正确的分流和指标设计。
  • 指标体系应分层:业务指标、产品指标、技术指标。
  • 数据看板需要实时性、准确性和可解释性。
  • 数据管道需要 Lambda 架构兼顾实时与离线场景。
  • 隐私合规已成为数据工程的硬性要求,必须从架构层面解决。

1. 数据埋点体系

1.1 埋点类型

类型说明示例
点击埋点用户点击元素按钮点击、链接点击
曝光埋点元素进入可视区广告曝光、商品曝光
页面埋点页面访问PV、UV、停留时长
自定义事件业务特定行为加入购物车、提交订单

1.2 埋点设计原则

  • 事件模型:Who(用户标识)、When(时间戳)、Where(页面/模块)、What(行为类型)、How(操作方式)。
  • 属性标准化:用户 ID、设备信息、页面路径、模块名称、元素标识、时间戳为通用属性。
  • 事件命名规范:采用 对象_动作 格式,如 button_clickpage_viewform_submit
  • 属性类型定义:字符串、数值、布尔值、数组、对象,确保数据类型一致。
  • 埋点规范文档化:建立埋点字典,包含事件名、属性列表、触发时机、责任人,避免口径不一致。

1.3 埋点 SDK 完整设计

以下是一个完整的前端埋点 SDK 实现,包含批量上报、重试机制、离线回退和自动采集。

javascript
/**
 * @typedef {Object} TrackingConfig
 * @property {string} endpoint - 上报接口地址
 * @property {number} [batchSize=10] - 批量上报最大条数
 * @property {number} [flushInterval=5000] - 自动刷新间隔(毫秒)
 * @property {number} [maxRetries=3] - 最大重试次数
 * @property {boolean} [autoPageView=true] - 是否自动采集页面浏览
 * @property {boolean} [autoClick=true] - 是否自动采集点击事件
 * @property {boolean} [debug=false] - 调试模式
 * @property {string} [storageKey='track_events'] - 本地存储键名
 */

/**
 * @typedef {Object} TrackEvent
 * @property {string} event - 事件名称
 * @property {Object} properties - 事件属性
 * @property {string} timestamp - 事件发生时间
 * @property {string} uuid - 事件唯一 ID
 * @property {number} retryCount - 已重试次数
 */

class TrackingSDK {
  /** @type {TrackingConfig} */
  config;

  /** @type {TrackEvent[]} */
  queue = [];

  /** @type {string} */
  anonymousId = '';

  /** @type {Object|null} */
  userInfo = null;

  /** @type {boolean} */
  isFlushing = false;

  /** @type {number|null} */
  flushTimer = null;

  /** @type {AbortController|null} */
  abortController = null;

  /** @type {boolean} */
  onlineStatus = true;

  constructor(config) {
    this.config = {
      batchSize: 10,
      flushInterval: 5000,
      maxRetries: 3,
      autoPageView: true,
      autoClick: true,
      debug: false,
      storageKey: 'track_events',
      ...config
    };
    this.anonymousId = this.generateId();
    this.restoreFromStorage();
    this.initAutoTracking();
    this.initNetworkListener();
    this.startFlushTimer();
  }

  /** 生成唯一标识 */
  generateId() {
    return 'xxxxxxxx-xxxx-4xxx-yxxx-xxxxxxxxxxxx'.replace(/[xy]/g, (c) => {
      const r = (Math.random() * 16) | 0;
      return (c === 'x' ? r : (r & 0x3) | 0x8).toString(16);
    });
  }

  /** 获取通用属性 */
  getCommonProperties() {
    return {
      anonymous_id: this.anonymousId,
      user_id: this.userInfo?.userId || null,
      url: window.location.href,
      referrer: document.referrer,
      user_agent: navigator.userAgent,
      screen: `${window.screen.width}x${window.screen.height}`,
      timestamp: new Date().toISOString(),
      timezone: Intl.DateTimeFormat().resolvedOptions().timeZone,
      language: navigator.language,
      platform: navigator.platform
    };
  }

  /**
   * 追踪事件
   * @param {string} event - 事件名称
   * @param {Object} [properties={}] - 事件属性
   */
  track(event, properties = {}) {
    if (this.config.debug) {
      console.log('[Tracking]', event, properties);
    }
    const trackEvent = {
      event,
      properties: { ...this.getCommonProperties(), ...properties },
      timestamp: new Date().toISOString(),
      uuid: this.generateId(),
      retryCount: 0
    };
    this.queue.push(trackEvent);
    if (this.queue.length >= this.config.batchSize) {
      this.flush();
    }
  }

  /**
   * 标识用户
   * @param {Object} user - 用户信息
   * @param {string} user.userId - 用户 ID
   * @param {Object} [user.traits={}] - 用户属性
   */
  identify(user) {
    this.userInfo = user;
    this.track('user_identify', {
      user_id: user.userId,
      ...(user.traits || {})
    });
  }

  /**
   * 页面浏览追踪
   * @param {string} [pageName] - 页面名称
   */
  page(pageName) {
    this.track('page_view', {
      page_name: pageName || document.title,
      page_path: window.location.pathname,
      page_title: document.title,
      search: window.location.search
    });
  }

  /**
   * 批量上报
   */
  async flush() {
    if (this.isFlushing || this.queue.length === 0) return;
    this.isFlushing = true;
    const batch = this.queue.splice(0, this.config.batchSize);
    try {
      await this.sendBatch(batch);
      if (this.config.debug) {
        console.log('[Tracking] Flushed', batch.length, 'events');
      }
    } catch (error) {
      console.error('[Tracking] Flush failed:', error);
      batch.forEach(e => { e.retryCount++ });
      const retryable = batch.filter(e => e.retryCount <= this.config.maxRetries);
      this.queue.unshift(...retryable);
      if (retryable.length < batch.length) {
        console.warn('[Tracking] Dropped', batch.length - retryable.length, 'events after max retries');
      }
    } finally {
      this.isFlushing = false;
      this.saveToStorage();
    }
  }

  /**
   * 发送批量事件到服务端
   * @param {TrackEvent[]} events
   */
  async sendBatch(events) {
    this.abortController = new AbortController();
    const timeout = setTimeout(() => this.abortController.abort(), 10000);
    try {
      const response = await fetch(this.config.endpoint, {
        method: 'POST',
        headers: { 'Content-Type': 'application/json' },
        body: JSON.stringify({ events }),
        signal: this.abortController.signal,
        keepalive: true
      });
      if (!response.ok) {
        throw new Error(`HTTP ${response.status}`);
      }
    } finally {
      clearTimeout(timeout);
      this.abortController = null;
    }
  }

  /** 保存到本地存储(离线回退) */
  saveToStorage() {
    try {
      const data = JSON.stringify(this.queue);
      localStorage.setItem(this.config.storageKey, data);
    } catch (e) {
      console.warn('[Tracking] Storage save failed:', e);
    }
  }

  /** 从本地存储恢复未上报事件 */
  restoreFromStorage() {
    try {
      const data = localStorage.getItem(this.config.storageKey);
      if (data) {
        const saved = JSON.parse(data);
        if (Array.isArray(saved)) {
          this.queue.unshift(...saved);
        }
        localStorage.removeItem(this.config.storageKey);
      }
    } catch (e) {
      console.warn('[Tracking] Storage restore failed:', e);
    }
  }

  /** 初始化自动追踪 */
  initAutoTracking() {
    if (this.config.autoPageView) {
      if (document.readyState === 'complete') {
        this.page();
      } else {
        window.addEventListener('load', () => this.page());
      }
      let lastUrl = location.href;
      new MutationObserver(() => {
        const url = location.href;
        if (url !== lastUrl) {
          lastUrl = url;
          this.page();
        }
      }).observe(document, { subtree: true, childList: true });
    }
    if (this.config.autoClick) {
      document.addEventListener('click', (e) => {
        const target = e.target;
        const trackData = target.getAttribute('data-track');
        if (trackData) {
          try {
            const parsed = JSON.parse(trackData);
            this.track('element_click', {
              element: target.tagName.toLowerCase(),
              text: target.innerText?.slice(0, 50),
              ...parsed
            });
          } catch {
            this.track('element_click', {
              element: target.tagName.toLowerCase(),
              text: target.innerText?.slice(0, 50),
              track_value: trackData
            });
          }
        }
      }, true);
    }
  }

  /** 初始化网络状态监听 */
  initNetworkListener() {
    window.addEventListener('online', () => {
      this.onlineStatus = true;
      this.flush();
    });
    window.addEventListener('offline', () => {
      this.onlineStatus = false;
    });
  }

  /** 启动定时刷新 */
  startFlushTimer() {
    this.flushTimer = setInterval(() => {
      if (this.queue.length > 0) {
        this.flush();
      }
    }, this.config.flushInterval);
  }

  /**
   * 销毁 SDK,清理资源
   */
  destroy() {
    if (this.flushTimer) {
      clearInterval(this.flushTimer);
      this.flushTimer = null;
    }
    if (this.abortController) {
      this.abortController.abort();
    }
    this.flush();
  }
}

// ======== Intersection Observer 曝光追踪 ========
class ImpressionTracker {
  /** @type {IntersectionObserver} */
  observer;

  /** @type {TrackingSDK} */
  sdk;

  constructor(sdk, options = {}) {
    this.sdk = sdk;
    this.observer = new IntersectionObserver(
      (entries) => {
        entries.forEach(entry => {
          if (entry.isIntersecting) {
            const el = entry.target;
            const trackData = el.getAttribute('data-impression');
            if (trackData) {
              try {
                const parsed = JSON.parse(trackData);
                sdk.track('element_impression', {
                  element: el.tagName.toLowerCase(),
                  ...parsed
                });
              } catch {
                sdk.track('element_impression', {
                  element: el.tagName.toLowerCase(),
                  impression_value: trackData
                });
              }
            }
            // 只曝光一次
            this.observer.unobserve(el);
          }
        });
      },
      { threshold: 0.5, ...options }
    );
  }

  /**
   * 观察元素曝光
   * @param {Element} el
   */
  observe(el) {
    this.observer.observe(el);
  }

  destroy() {
    this.observer.disconnect();
  }
}

// ======== 使用示例 ========
const tracker = new TrackingSDK({
  endpoint: 'https://api.example.com/track',
  batchSize: 20,
  flushInterval: 10000,
  debug: true
});

tracker.identify({
  userId: 'user_12345',
  traits: { plan: 'premium', age: 28 }
});

tracker.track('button_click', {
  button_name: 'signup_submit',
  page: 'landing'
});

const impressionTracker = new ImpressionTracker(tracker);
document.querySelectorAll('[data-impression]').forEach(el => {
  impressionTracker.observe(el);
});

1.4 Data Layer 与事件模型设计

数据层是整个埋点体系的核心,需要统一管理公共属性、事件属性和用户属性。

属性类别字段说明示例
公共属性app_id应用 ID'web_main'
公共属性platform平台标识'web', 'mini_app'
公共属性versionSDK 版本'1.2.0'
公共属性locale用户语言区域'zh-CN'
事件属性event_name事件名称'button_click'
事件属性event_id事件唯一 IDUUID v4
事件属性properties自定义属性任意 JSON
用户属性user_id登录用户 ID'u_10086'
用户属性device_id设备指纹SHA256 哈希
用户属性session_id会话 ID每次访问生成

数据质量校验管道

javascript
function validateEvent(event) {
  const errors = [];
  if (!event.event || typeof event.event !== 'string') {
    errors.push('event: 必须是非空字符串');
  }
  if (!/^[a-z][a-z0-9_]*$/.test(event.event)) {
    errors.push('event: 必须是小写字母开头,仅含小写字母、数字和下划线');
  }
  if (!event.timestamp || isNaN(Date.parse(event.timestamp))) {
    errors.push('timestamp: 必须为有效 ISO 日期字符串');
  }
  if (!event.uuid || typeof event.uuid !== 'string') {
    errors.push('uuid: 必须为非空字符串');
  }
  const commonFields = ['url', 'anonymous_id'];
  for (const field of commonFields) {
    if (!event.properties?.[field]) {
      errors.push(`properties.${field}: 缺少必要公共属性`);
    }
  }
  return {
    valid: errors.length === 0,
    errors
  };
}

1.5 Schema 管理与版本控制

  • Schema Registry:集中管理所有事件的 Schema 定义,支持 JSON Schema 格式。
  • 版本演进:事件属性只能新增不能删除,标记为 deprecated 而非直接移除。
  • 兼容性检查:每次部署前自动检查 Schema 向后兼容性。
  • 自动化测试:埋点验收测试(SAT, Source Acceptance Test)确保埋点代码与 Schema 一致。
  • PII 脱敏:敏感字段(手机号、邮箱、身份证)在 SDK 层自动脱敏,不上报原始值。
javascript
/** PII 数据脱敏函数 */
function maskPII(value, type) {
  if (!value) return value;
  const str = String(value);
  switch (type) {
    case 'phone':
      return str.replace(/(\\d{3})\\d{4}(\\d{4})/, '$1****$2');
    case 'email':
      const [name, domain] = str.split('@');
      return name[0] + '***@' + domain;
    case 'id_card':
      return str.slice(0, 4) + '***********' + str.slice(-4);
    default:
      return '***MASKED***';
  }
}

1.6 Debug 模式与开发体验

  • Debug Mode:在 SDK 初始化时开启 debug: true,在控制台输出所有事件的详细信息。
  • 可视化排查工具:浏览器扩展或嵌入面板,实时展示已采集的事件队列。
  • 事件回放:通过 DevTools 的 Network 面板查看上报请求的 payload。
  • Mock 模式:在测试环境中将 endpoint 替换为本地 Mock 服务器,不产生真实流量。
  • 日志分级debug(全量)、info(核心事件)、warn(校验警告)、error(上报失败)。

2. AB 测试平台架构

2.1 实验流程总览

AB 测试是验证产品假设的黄金标准方法。完整的实验流程如下:

提出假设 → 设计实验 → 随机分流 → 运行实验 → 收集数据 → 统计分析 → 决策判断

2.2 平台架构组件

一个完整的 AB 测试平台包含以下核心服务:

服务职责技术栈示例
Experiment Service实验配置管理、生命周期关系数据库 + Admin UI
Assignment Service用户分流决策、实时低延迟Redis + Hash 算法
Metrics Service指标定义与数据聚合ClickHouse + Presto
Analysis Service统计分析、P 值计算Python (statsmodels) + R
Feature Flag Service功能开关集成LaunchDarkly / 自建

2.3 分流算法实现

分流的核心在于确定性:同一用户在同一个实验中始终看到同一版本。

javascript
/**
 * Hash 分流算法
 * @param {string} userId - 用户 ID
 * @param {string} experimentId - 实验 ID
 * @param {number} totalVariants - 分桶总数
 * @returns {number} 分桶下标(从 0 开始)
 */
function hashAssign(userId, experimentId, totalVariants) {
  const key = `${experimentId}:${userId}`;
  let hash = 0;
  for (let i = 0; i < key.length; i++) {
    const char = key.charCodeAt(i);
    hash = ((hash << 5) - hash) + char;
    hash = hash & hash;
  }
  return Math.abs(hash) % totalVariants;
}

/**
 * 流量分配:给每个版本分配百分比
 * @param {string} userId - 用户 ID
 * @param {string} experimentId - 实验 ID
 * @param {number[]} trafficAllocations - 各版本流量占比,如 [0.5, 0.5]
 * @returns {number} 版本下标
 */
function trafficSplitAssign(userId, experimentId, trafficAllocations) {
  const bucket = hashAssign(userId, experimentId, 10000);
  let cumulative = 0;
  for (let i = 0; i < trafficAllocations.length; i++) {
    cumulative += trafficAllocations[i] * 10000;
    if (bucket < cumulative) return i;
  }
  return trafficAllocations.length - 1;
}

多层分流策略对比

策略粒度优势劣势
用户级分流用户体验一致,适合 UI 实验样本量要求大
设备级分流设备跨设备一致性设备 ID 获取复杂
会话级分流会话样本量大体验可能不一致
页面级分流曝光最大样本量用户感知明显

2.4 统计分析方法

频率派 vs 贝叶斯派

方法频率派(Frequentist)贝叶斯派(Bayesian)
核心P 值假设检验后验概率分布
停止规则固定样本量可随时停止(连续监测)
解释如果无效假设为真,观察到的结果概率实验组优于对照组的概率
工具t-test, Z-test, chi-squareBeta-Binomial, MCMC
多重比较修正Bonferroni, Benjamini-Hochberg分层贝叶斯模型

独立双样本 t 检验实现

javascript
/**
 * 独立双样本 t 检验计算
 * @param {number[]} control - 对照组数据
 * @param {number[]} treatment - 实验组数据
 * @returns &#123;&#123; t: number, p: number, significant: boolean &#125;&#125;
 */
function tTest(control, treatment) {
  const n1 = control.length, n2 = treatment.length;
  const mean1 = control.reduce((a, b) => a + b, 0) / n1;
  const mean2 = treatment.reduce((a, b) => a + b, 0) / n2;
  const var1 = control.reduce((s, v) => s + (v - mean1) ** 2, 0) / (n1 - 1);
  const var2 = treatment.reduce((s, v) => s + (v - mean2) ** 2, 0) / (n2 - 1);
  const se = Math.sqrt(var1 / n1 + var2 / n2);
  const t = (mean1 - mean2) / se;
  const df = (var1 / n1 + var2 / n2) ** 2 /
    ((var1 / n1) ** 2 / (n1 - 1) + (var2 / n2) ** 2 / (n2 - 1));
  // 使用 t 分布近似计算 P 值(简化版)
  const p = 2 * (1 - studentT_CDF(Math.abs(t), df));
  return { t, p, significant: p < 0.05 };
}

function studentT_CDF(x, df) {
  // 使用正则化不完全 Beta 函数近似
  const t = df / (df + x * x);
  return 1 - 0.5 * ibeta(t, df / 2, 0.5);
}

2.5 统计功效与样本量计算

在实验设计阶段,需要预先计算所需样本量,确保实验有足够的统计功效(通常要求 80%)。

javascript
/**
 * 最小样本量估算
 * @param {number} baseline - 基线转化率
 * @param {number} minDetectableEffect - 最小可检测提升(如 0.05 表示 5%)
 * @param {number} alpha - 显著性水平(默认 0.05)
 * @param {number} power - 统计功效(默认 0.80)
 * @returns {number} 每个版本所需最小样本量
 */
function minSampleSize(baseline, minDetectableEffect, alpha = 0.05, power = 0.80) {
  const zAlpha = 1.96;
  const zPower = 0.84;
  const p1 = baseline;
  const p2 = baseline * (1 + minDetectableEffect);
  const pBar = (p1 + p2) / 2;
  const numerator = (zAlpha * Math.sqrt(2 * pBar * (1 - pBar)) +
    zPower * Math.sqrt(p1 * (1 - p1) + p2 * (1 - p2))) ** 2;
  const denominator = (p2 - p1) ** 2;
  return Math.ceil(numerator / denominator);
}

2.6 样本比率不匹配(SRM)检测

SRM(Sample Ratio Mismatch)是 AB 实验中最常见也最危险的问题之一。当实际分流比例与预期比例出现显著偏差时,说明分流可能被污染。

javascript
/**
 * 卡方检验检测 SRM
 * @param {number[]} observed - 观察到的各版本样本量
 * @param {number[]} expected - 预期的各版本样本量
 * @returns &#123;&#123; chi2: number, p: number, hasSRM: boolean &#125;&#125;
 */
function detectSRM(observed, expected) {
  const totalObs = observed.reduce((a, b) => a + b, 0);
  const totalExp = expected.reduce((a, b) => a + b, 0);
  let chi2 = 0;
  for (let i = 0; i < observed.length; i++) {
    const expectedCount = totalObs * (expected[i] / totalExp);
    chi2 += (observed[i] - expectedCount) ** 2 / expectedCount;
  }
  const df = observed.length - 1;
  const p = 1 - chiSquaredCDF(chi2, df);
  return { chi2, p, hasSRM: p < 0.05 };
}

2.7 多重比较修正

当同时运行多个实验或多个指标时,假阳性率会显著增加。修正方法如下:

方法说明控制指标保守程度
Bonferroni将阈值除以比较次数FWER最保守
Holm-Bonferroni逐步向下检验FWER中等
Benjamini-Hochberg控制错误发现率FDR较不保守
Storey 方法基于估计的 FDRFDR最不保守

2.8 实验生命周期管理

每个实验经过以下阶段:

阶段操作数据要求
Draft创建实验、定义指标、设置分流
Running开始实验、分配流量、收集数据实时监控护栏指标
Analysis停止入流、数据已全部到达完整数据 + 统计检验
Decision判断是否显著、做出决策分析报告
Cleanup清理实验代码、归档结果最终报告存档

3. 指标树设计(Metric Tree)

3.1 指标分层模型

指标树是连接公司战略与具体执行的核心工具。标准分层模型如下:

北极星指标(North Star Metric)
    └── 目标指标(Goal Metrics)
           └── 输入指标(Input Metrics)
                  └── 驱动指标(Driver Metrics)
层级定义示例(电商)示例(社区)
北极星指标公司级核心成功指标GMVDAU
目标指标子业务线成功指标下单转化率发帖率
输入指标团队可控的输入指标加购率评论率
驱动指标具体可执行的指标商品详情页加载时间照片上传成功率

3.2 Goal-Signal-Metric 框架

每个指标应从三个维度明确定义:

  • Goal(目标):该指标对应的业务目标是什幺?
  • Signal(信号):通过什么数据行为来反映目标达成?
  • Metric(指标):如何量化该信号?

示例——提升搜索质量:

  • Goal:用户更快找到想要的商品
  • Signal:用户搜索后点击第一个结果
  • Metric:搜索结果 Top-1 点击率(CTR)

3.3 指标定义模板

字段说明示例值
指标名称指标的唯一标识search_ctr
指标分类所属层级和类别产品指标 / 搜索
指标定义计算公式描述搜索后点击结果数 / 搜索总次数
计算公式精确的 SQL 表达式SUM(CASE WHEN click>0 THEN 1 ELSE 0 END) / COUNT(*)
数据来源采集数据的系统搜索日志埋点
采集频率数据更新频率实时(T+0)
目标值期望达成的数值> 0.45
报警阈值触发告警的阈值< 0.35
责任人指标负责团队搜索产品组
北极星关联与北极星的关系正相关

3.4 计数器指标与比率指标

类型定义示例注意点
计数器直接计数PV、UV、订单数易受总量影响,需要归一化
比率指标分子/分母CTR、转化率分母为 0 的处理很关键
均值指标总和/计数平均停留时长异常值敏感
分位数指标中位数 / P95P95 加载时间抗异常值,但计算成本高

3.5 滞后指标 vs 先行指标

  • 滞后指标(Lagging Indicators):反映已经发生的结果,如 GMV、留存率。适合回顾性分析,不适合日常优化。
  • 先行指标(Leading Indicators):预测未来结果的指标,如活跃用户数、功能采用率。可用于日常决策和快速迭代。
  • 策略:每个滞后指标至少配套 2-3 个先行指标,形成完整的监控闭环。

3.6 不同产品类型的北极星指标

产品类型北极星指标理由
电商平台有效 GMV直接反映交易规模和健康度
社交网络MAU(月活)网络效应与用户粘性的核心
内容平台阅读时长用户注意力是最核心的资产
SaaS 产品付费客户数持续付费验证产品价值
工具产品任务完成率工具的本质是完成任务
游戏付费用户平均收入(ARPPU)收入与体验的平衡

4. 数据管道(Data Pipeline)

4.1 Lambda 架构总览

现代数据管道通常采用 Lambda 架构,同时处理实时和离线数据。

                   ┌─────────────┐
                   │  用户行为      │
                   │  埋点数据      │
                   └──────┬──────┘

              ┌───────────┴───────────┐
              │                       │
      ┌───────▼───────┐     ┌────────▼───────┐
      │  实时管道       │     │  离线管道       │
      │  Kafka → Flink │     │  HDFS → Spark  │
      │  → Redis       │     │  → Hive        │
      └───────┬───────┘     └────────┬───────┘
              │                       │
              └───────────┬───────────┘

                   ┌──────▼──────┐
                   │  OLAP 存储   │
                   │  ClickHouse  │
                   │  / Druid     │
                   └──────┬──────┘

                   ┌──────▼──────┐
                   │  数据看板    │
                   │  Superset   │
                   │  / Grafana  │
                   └─────────────┘

4.2 采集层(Collection Layer)

  • SDK:前端 SDK 采集事件,通过 navigator.sendBeacon() 或 fetch API 上报。
  • API Gateway:统一的接收端点,负责鉴权、限流、数据格式校验。
  • Kafka:事件流的缓冲和分发层,保证数据不丢失。
  • Schema Registry:事件 Schema 管理,兼容性校验,序列化/反序列化。

4.3 ETL 层(Extract, Transform, Load)

javascript
/**
 * ETL 管道中的事件清洗函数示例
 */
function cleanEvent(raw) {
  // 1. 提取必需字段
  const cleaned = {
    event: raw.event?.trim() || 'unknown',
    properties: {}, 
    timestamp: raw.timestamp || new Date().toISOString(),
    uuid: raw.uuid || crypto.randomUUID()
  };
  // 2. 类型转换(数字字段转换)
  if (raw.properties?.price) {
    cleaned.properties.price = Number(raw.properties.price);
  }
  // 3. 字段标准化(统一大小写)
  for (const [key, value] of Object.entries(raw.properties || {})) {
    cleaned.properties[key.toLowerCase()] = value;
  }
  // 4. 非法字符过滤
  cleaned.event = cleaned.event.replace(/[^a-z0-9_]/g, '');
  // 5. 衍生字段(补充时间维度)
  const dt = new Date(cleaned.timestamp);
  cleaned.properties._date = dt.toISOString().slice(0, 10);
  cleaned.properties._hour = dt.getUTCHours();
  cleaned.properties._weekday = dt.getUTCDay();
  return cleaned;
}

ETL 的四个核心步骤:

  • Validation(校验):格式校验、必填字段校验、Schema 兼容性校验。
  • Cleaning(清洗):去除重复事件、修正格式错误、过滤爬虫机器人流量。
  • Enrichment(丰富):补充 IP 地理位置、UA 解析设备信息、用户画像标签。
  • Transformation(转换):格式转换、单位统一、派生字段计算。

4.4 数仓建模(Warehouse Layer)

星型模型 vs 雪花模型

特征星型模型(Star Schema)雪花模型(Snowflake Schema)
事实表中央事实表,包含度量值和维度外键同上
维度表非规范化,冗余存储规范化,多层级关联
查询性能更优(JOIN 少)较差(JOIN 多)
存储空间更多(冗余)更少(无冗余)
维护成本较低较高

典型事实表设计

sql
-- 事件事实表
CREATE TABLE events_fact (
  event_id VARCHAR(64),
  event_name VARCHAR(128),
  user_id VARCHAR(64),
  device_id VARCHAR(64),
  session_id VARCHAR(64),
  page_id INTEGER,
  timestamp DATETIME,
  date_key INTEGER,  -- YYYYMMDD 格式
  location_id INTEGER,
  properties JSON,
  PRIMARY KEY (event_id)
);

-- 维度表:日期
CREATE TABLE dim_date (
  date_key INTEGER PRIMARY KEY,
  full_date DATE,
  year INTEGER,
  month INTEGER,
  day INTEGER,
  weekday INTEGER,
  is_holiday BOOLEAN
);

5. ClickHouse 实时分析

5.1 为什么选择 ClickHouse

ClickHouse 是专为 OLAP 场景设计的列式存储数据库,在前端数据工程中有以下优势:

  • 列式存储:只读取需要的列,IO 效率极高。
  • 向量化计算:利用 CPU SIMD 指令并行处理数据块。
  • 实时写入:支持每秒百万行级别的写入吞吐。
  • 极速查询:典型聚合查询在毫秒到秒级完成。
  • 数据压缩:列式存储配合 LZ4/ZSTD 压缩算法,压缩比可达 5-10 倍。

5.2 表引擎家族

引擎适用场景特点
MergeTree通用事件数据主键排序、分区、TTL
ReplacingMergeTree需要去重的数据按排序键去重(最终一致性)
AggregatingMergeTree预聚合指标增量合并聚合状态
SummingMergeTree计数类指标自动合并 SUM 结果
Distributed分布式查询透明访问集群中所有分片

事件表设计

sql
CREATE TABLE events_local (
  event_date Date,
  event_time DateTime,
  event_name String,
  user_id String,
  session_id String,
  page_url String,
  referrer String,
  device_type String,
  country String,
  properties String  -- JSON 字符串
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_name, cityHash64(user_id))
TTL event_date + INTERVAL 90 DAY DELETE;

-- 分布式表
CREATE TABLE events_all AS events_local
ENGINE = Distributed('cluster_name', 'default', 'events_local', cityHash64(user_id));

物化视图(Materialized Views)

sql
-- 每小时 PV/UV 统计物化视图
CREATE MATERIALIZED VIEW hourly_metrics_mv
ENGINE = AggregatingMergeTree()
PARTITION BY toYYYYMM(hour)
ORDER BY (hour, event_name)
AS SELECT
  toStartOfHour(event_time) AS hour,
  event_name,
  count() AS pv,
  uniqState(user_id) AS uv_state
FROM events_all
GROUP BY hour, event_name;

-- 查询物化视图
SELECT
  hour,
  event_name,
  pv,
  uniqMerge(uv_state) AS uv
FROM hourly_metrics_mv
WHERE hour >= now() - INTERVAL 7 DAY
GROUP BY hour, event_name
ORDER BY hour DESC;

5.3 漏斗分析 SQL

sql
-- 电商漏斗:浏览 → 加购 → 下单 → 支付
WITH
funnel AS (
  SELECT
    user_id,
    maxIf(event_name = 'page_view_product', 1, 0) AS step1,
    maxIf(event_name = 'add_to_cart', 1, 0) AS step2,
    maxIf(event_name = 'create_order', 1, 0) AS step3,
    maxIf(event_name = 'payment_success', 1, 0) AS step4
  FROM events_all
  WHERE event_date >= '2026-01-01' AND event_date < '2026-02-01'
  GROUP BY user_id
)
SELECT
  '浏览商品' AS step, countIf(step1 = 1) AS users, 100.0 AS pct
FROM funnel
UNION ALL
SELECT
  '加入购物车' AS step, countIf(step2 = 1) AS users,
  round(countIf(step2 = 1) * 100.0 / countIf(step1 = 1), 1) AS pct
FROM funnel
UNION ALL
SELECT
  '创建订单' AS step, countIf(step3 = 1) AS users,
  round(countIf(step3 = 1) * 100.0 / countIf(step2 = 1), 1) AS pct
FROM funnel
UNION ALL
SELECT
  '支付成功' AS step, countIf(step4 = 1) AS users,
  round(countIf(step4 = 1) * 100.0 / countIf(step3 = 1), 1) AS pct
FROM funnel;

5.4 留存分析 SQL

sql
-- 日留存分析
WITH
first_visit AS (
  SELECT user_id, toDate(min(event_time)) AS first_date
  FROM events_all
  WHERE event_name = 'page_view'
  GROUP BY user_id
),
daily_visits AS (
  SELECT DISTINCT user_id, toDate(event_time) AS visit_date
  FROM events_all
  WHERE event_name = 'page_view'
)
SELECT
  fv.first_date,
  count(DISTINCT fv.user_id) AS new_users,
  countIf(dv.visit_date = fv.first_date + 1) AS day1,
  countIf(dv.visit_date = fv.first_date + 7) AS day7,
  countIf(dv.visit_date = fv.first_date + 30) AS day30
FROM first_visit fv
LEFT JOIN daily_visits dv ON fv.user_id = dv.user_id
GROUP BY fv.first_date
ORDER BY fv.first_date;

5.5 ClickHouse vs 传统 OLAP 对比

维度ClickHouseMySQLPresto/Trino
存储模型列式行式列式(内存)
写入性能百万行/秒万行/秒N/A
聚合查询毫秒级秒到分钟级秒级
JOIN 性能良好(大表 LEFT JOIN)一般优秀
并发中(100+ 查询)
数据压缩5-10x2-3x无固定
适用场景实时 OLAPOLTP交互式查询

6. 隐私合规(Privacy Compliance)

6.1 全球隐私法规概览

法规区域核心要求处罚力度
GDPR欧盟/EEA同意获取、数据可删除、可携带年营收 4% 或 2000 万欧元
CCPA/CPRA美国加州知情权、删除权、拒绝出售每条违规 2500-7500 美元
PIPL中国个人信息处理规则、最小必要最高 5000 万或年营收 5%
LGPD巴西类似 GDPR年营收 2% 最高 5000 万雷亚尔

6.2 同意管理平台(CMP)集成

javascript
/**
 * 与 CMP(Consent Management Platform)集成
 * 在用户授权后方可开始采集
 */
class ConsentManager {
  /** @type {'granted'|'denied'|'pending'} */
  consentStatus = 'pending';

  /** @type {Set<string>} */
  consentedPurposes = new Set();

  async init() {
    // 等待 CMP SDK 加载完成
    await new Promise(resolve => {
      if (window.__cmp) { resolve(); }
      else { window.addEventListener('cmp-ready', resolve, { once: true }); }
    });
    const data = await this.getConsentData();
    this.consentStatus = data.granted ? 'granted' : 'denied';
    this.consentedPurposes = new Set(data.purposes);
    return this.consentStatus;
  }

  canTrack(purpose) {
    return this.consentStatus === 'granted' &&
      (!purpose || this.consentedPurposes.has(purpose));
  }
}

随着第三方 Cookie 逐步淘汰,需要采用替代方案:

方案原理优点局限
First-Party Cookie使用自身域名的 Cookie浏览器支持良好跨域无法共享
LocalStorage浏览器本地存储API 简单,容量大可手动清除
设备指纹Canvas/WebGL/Font 指纹无需存储识别率非 100%,隐私争议
Server-Side Cookie服务端设置 Cookie更稳定需要自有域名
Private Click Measurement浏览器 API(如 Safari PCM)隐私保护功能有限,数据延迟

6.4 数据保留策略

数据类型保留期限理由
原始事件数据90 天分析实时趋势、短期指标
聚合指标数据永久长期趋势和 YoY 对比
用户 PII 数据账号注销后 30 天删除GDPR 数据最小化原则
AB 实验原始数据实验结束后 90 天实验复验和审计
日志备份180 天安全审计和故障排查

6.5 匿名化 vs 假名化

方法定义可逆性GDPR 适用
匿名化(Anonymization)彻底去除所有可识别信息不可逆不适用 GDPR
假名化(Pseudonymization)用假名标识符替代直接标识可逆(通过映射表)仍适用 GDPR
聚合(Aggregation)只保留统计结果不可逆不适用 GDPR

隐私优先的 SDK 设计原则

  • 默认不采集任何数据(Opt-in)。
  • 支持按目的分类授权(功能必要、分析、个性化、广告)。
  • SDK 层在采集前检查授权状态。
  • 敏感字段在客户端脱敏。
  • 提供数据下载和删除接口(GDPR 数据主体权利)。

7. 漏斗 / 留存 / 同期群分析(Funnel / Cohort / Retention)

7.1 漏斗分析

漏斗分析用于追踪用户在多步骤流程中的转化和流失情况。

步骤用户数转化率流失率
访问首页100,000100.0%0.0%
搜索商品65,00065.0%35.0%
查看商品详情32,50050.0%50.0%
加入购物车13,00040.0%60.0%
完成支付5,20040.0%60.0%

关键分析维度

  • 步骤间流失:哪一步流失最大?这是优化重点。
  • 分群对比:不同渠道、设备、地区的转化差异。
  • 时间趋势:漏斗转化率随时间的变化趋势。
  • 反向漏斗:从终点倒推,分析成功用户的行为路径。

7.2 留存分析

留存分析衡量用户在一段时间后是否仍在使用产品。

指标定义说明
Day N 留存第 N 天回访用户 / 首日用户产品粘性的核心指标
Bracket 留存第 N 周/月回访适用于低频产品
Unbounded 留存任一天回访宽松指标
滚动留存N 天内至少回访一次适合内容产品

7.3 同期群分析(Cohort Analysis)

同期群分析将用户按首次行为时间分组,追踪各组在后续周期中的表现。

同期群留存表(示例):

         Week1  Week2  Week3  Week4  Week5
Cohort1  100%   45%    32%    28%    25%
Cohort2  100%   48%    35%    30%    -
Cohort3  100%   42%    30%    -      -
Cohort4  100%   50%    -      -      -
Cohort5  100%   -      -      -      -

同期群类型

  • 时间同期群:按周/月注册时间分组。
  • 行为同期群:按首次行为类型分组(首次购买、首次发帖)。
  • 渠道同期群:按获客渠道分组(自然流量、付费广告、社交推荐)。
  • 规模同期群:按用户消费金额或频次分组。

7.4 用户留存曲线与幂律分布

用户留存曲线通常遵循幂律分布(Power Law):

  • 前 7 天:留存率快速下降,这是用户习惯养成的关键窗口。
  • 7-30 天:下降速度放缓,核心用户群体逐渐稳定。
  • 30-90 天:趋于平稳,这是产品的长期留存基线。
  • 策略:通过新用户引导(Onboarding)提升首周留存;通过产品价值强化提升长期留存。

8. 实时看板架构(Real-time Dashboard)

8.1 实时数据管道架构

埋点 SDK → API Gateway → Kafka → Flink → Redis → WebSocket Server → 前端看板

8.2 各组件职责

组件技术职责
数据采集SDK + API Gateway接收原始事件,格式校验
消息队列Kafka事件缓冲、削峰填谷、多消费者分发
流处理Flink / Kafka Streams实时 ETL、窗口聚合、维度关联
状态存储Redis实时计数器、HyperLogLog、排行榜
消息推送WebSocket / SSE将聚合结果推送到前端
前端展示React / Vue + ECharts实时图表更新

8.3 增量聚合与 HyperLogLog

对于 UV 等去重指标,使用 HyperLogLog(HLL)进行近似去重,兼顾精度和性能。

javascript
// 模拟 HLL 实时 UV 聚合
class RealtimeUV {
  constructor(redis) {
    this.redis = redis;
  }

  async track(pageId, userId) {
    const key = `uv:${pageId}:${this.getCurrentMinute()}`;
    await this.redis.pfadd(key, userId);
    // 设置过期时间,自动清理旧数据
    await this.redis.expire(key, 3600);
  }

  async getUV(pageId, minutes = 5) {
    const keys = [];
    for (let i = 0; i < minutes; i++) {
      const key = `uv:${pageId}:${this.getCurrentMinute(i)}`;
      keys.push(key);
    }
    return await this.redis.pfcount(...keys);
  }

  getCurrentMinute(offset = 0) {
    const d = new Date(Date.now() - offset * 60000);
    return `${d.getUTCFullYear()}${String(d.getUTCMonth()+1).padStart(2,'0')}${String(d.getUTCDate()).padStart(2,'0')}${String(d.getUTCHours()).padStart(2,'0')}${String(d.getMinutes()).padStart(2,'0')}`;
  }
}

8.4 看板组件架构

javascript
// React 实时看板组件示例
function RealtimeDashboard({ metrics }) {
  const [data, setData] = useState(null);

  useEffect(() => {
    const ws = new WebSocket('wss://dashboard.example.com/realtime');
    ws.onmessage = (event) => {
      const payload = JSON.parse(event.data);
      setData(payload);
    };
    return () => ws.close();
  }, []);

  if (!data) return <div>加载中...</div>;

  return (
    <div className='dashboard-grid'>
      <KpiCard title='实时 PV' value={data.pv} />
      <KpiCard title='实时 UV' value={data.uv} />
      <KpiCard title='平均停留' value={data.avgDuration + 's'} />
      <KpiCard title='错误率' value={data.errorRate + '%'} />
      <TimeseriesChart data={data.timeseries} />
      <TopPagesTable data={data.topPages} />
    </div>
  );
}

8.5 实时 vs 近实时权衡

场景延迟要求架构选择
监控告警秒级实时流处理(Flink)
产品看板分钟级近实时(微批次)
财务指标T+1离线批处理
个性化推荐毫秒级实时特征 + 预计算

看板性能优化策略

  • 预聚合:使用物化视图提前计算常用指标。
  • 缓存分层:浏览器缓存 → CDN → Redis → 数据库。
  • 查询超时:设置合理的查询超时,避免慢查询阻塞。
  • 数据采样:超大表使用抽样查询(SAMPLE)。
  • 增量加载:新数据增量追加,避免全表重算。

9. 特征工程 for ML

9.1 从原始事件到特征

用户行为数据经过特征工程可以转化为机器学习模型的输入。特征分为三大类:

特征类别示例提取方式
用户特征年龄、性别、会员等级用户画像表
上下文特征时间、设备、地理位置事件公共属性
行为特征过去 7 天点击次数、购买金额事件聚合

9.2 时间窗口特征提取

行为特征通常基于滑动时间窗口计算。

sql
-- ClickHouse 中提取用户 7 天特征
SELECT
  user_id,
  countIf(event_date >= today() - 7 AND event_name = 'page_view') AS view_count_7d,
  countIf(event_date >= today() - 7 AND event_name = 'add_to_cart') AS cart_count_7d,
  countIf(event_date >= today() - 7 AND event_name = 'payment_success') AS purchase_count_7d,
  countIf(event_date >= today() - 1 AND event_name = 'page_view') AS view_count_1d,
  uniqIf(page_url, event_date >= today() - 7) AS unique_pages_7d
FROM events_all
WHERE event_date >= today() - 7
GROUP BY user_id;

9.3 特征存储与在线服务

javascript
class FeatureStore {
  constructor(redis) {
    this.redis = redis;
  }

  /**
   * 计算并缓存用户特征
   * @param {string} userId
   */
  async computeFeatures(userId) {
    const cacheKey = `features:${userId}`;
    // 尝试从缓存获取
    const cached = await this.redis.get(cacheKey);
    if (cached) return JSON.parse(cached);
    // 从 ClickHouse 提取特征
    const features = await this.queryFeatures(userId);
    // 缓存 1 小时
    await this.redis.setex(cacheKey, 3600, JSON.stringify(features));
    return features;
  }
}

9.4 特征变换方法

方法适用于公式
标准化(Z-score)正态分布数据(x - mean) / std
归一化(Min-Max)有界区间数据(x - min) / (max - min)
对数变换长尾分布数据log(1 + x)
分箱(Binning)连续值离散化等宽 / 等频 / 分位数
One-Hot 编码无序类别每个类别一列
目标编码(Target Encoding)高基类别类别均值替代

9.5 特征验证与监控

  • 完整性检查:特征缺失率不能超过阈值(如 5%)。
  • 分布监控:特征分布漂移检测(PSI / KS 检验)。
  • 相关性分析:特征之间相关性过高需要合并或剔除。
  • 特征重要性:通过 SHAP / Permutation Importance 评估特征贡献。
  • 在线验证:A/B 测试对比使用/不使用特征的模型效果。

9.6 特征重要性分析示例

javascript
/**
 * 使用随机森林评估特征重要性(简化伪代码)
 */
function featureImportance(features, labels) {
  // 1. 训练随机森林模型
  const model = trainRandomForest(features, labels);
  // 2. 计算每个特征的基尼重要性
  const importances = model.featureImportances();
  // 3. 排序并返回
  return features.columns
    .map((name, i) => ({ name, importance: importances[i] }))
    .sort((a, b) => b.importance - a.importance);
}

10. 最佳实践与常见陷阱

10.1 埋点治理

  • 埋点即代码:埋点代码应纳入 Code Review 流程,与业务代码同等管理。
  • 验收测试:每次埋点上线前,通过自动化测试验证事件名称、属性、触发时机符合预期。
  • 数据质量监控:建立数据质量看板,监控事件量异常、属性缺失率、Schema 违规率。
  • 埋点生命周期:废弃的埋点及时下线,避免无效数据污染。
  • 指标口径字典:统一维护指标定义文档,避免不同团队口径不一致。

10.2 AB 实验避坑指南

陷阱后果解决方案
样本污染对照组受实验组影响隔离实验组和用户群体
新奇效应短期效果被高估延长实验周期,对比长期效果
实验时间不足统计功效不够提前计算所需样本量
多重比较假阳性率增加Bonferroni / FDR 修正
P-hacking选择性报告显著结果预先注册实验和指标
互斥实验实验间相互干扰建立实验层和互斥组

10.3 数据驱动文化

  • 数据可信:数据质量是第一位的,脏数据比没有数据更危险。
  • 指标可解释:团队每个成员都应该理解核心指标的含义和计算方式。
  • 假设先行:先提出可验证的假设,再通过数据验证,避免“先看数据再说”。
  • 闭环反馈:数据的最终目的是驱动决策,数据 → 洞察 → 行动 → 验证形成闭环。
  • 实验文化:鼓励尝试,容忍失败,但每次实验必须有明确的决策标准和止损机制。

常见误区

误区正确理解
埋点越多越好埋点要围绕核心指标,过多会增加治理成本
AB 实验只看核心指标还要关注护栏指标和长期影响
数据看板等于监控看板面向决策,监控面向告警
忽略数据质量脏数据会导致错误决策
离线和实时用同一管道两者有不同的 SLA 和技术要求
隐私合规是法务的事需要从架构层面落实 Privacy by Design

相关领域

  • L01 Business:业务理解和指标设计。
  • A06 Observability:技术可观测性,与前端监控互补。
  • E09 AI Engineering:数据驱动的智能应用。
  • A10 Visualization:数据可视化技术。
  • A11 Performance:性能优化,数据采集不能影响页面性能。
  • S01 Security:数据安全与隐私保护。

标签#data-engineering #ab-testing #metrics #dashboard #tracking

最后更新:2026-07-06


本领域学习进度

学习进度0 / 43 (0%)

基于 MIT 协议发布