Skip to content

大数据与商业分析资源专题

Big Data & Business Analytics Resources


🧭 学习路径

阶段 核心技术 工具
🟢 基础 SQL、数据仓库、ETL PostgreSQL、dbt
🟡 大数据 分布式计算、列式存储 Spark、Parquet、Hive
🟠 流处理 实时管道、事件驱动 Kafka、Flink
🔴 平台 数据湖、数据网格 Databricks、Snowflake

数据工程全貌 → 数据工程资源


📚 核心领域

大数据技术栈

技术领域 工具/框架 用途
分布式存储 Hadoop HDFS 海量数据存储
分布式计算 Apache Spark 内存计算引擎
数据仓库 Hive, Presto SQL 查询分析
流处理 Flink, Kafka Streams 实时数据处理
消息队列 Apache Kafka 事件流平台
NoSQL MongoDB, Cassandra 非结构化数据存储

商业分析

核心概念

概念 说明 应用场景
OMTM 唯一重要指标 初创公司聚焦
North Star Metric 北极星指标 产品方向指引
AARRR 海盗指标 用户增长模型
LTV 用户生命周期价值 用户价值评估
CAC 获客成本 营销效率
Churn Rate 流失率 用户留存
ARPU 每用户平均收入 收入分析

🎓 在线课程

大数据

课程 平台 链接
Big Data Specialization Coursera 链接
Spark 实战 Udemy 链接
Kafka 权威指南 O'Reilly 链接

商业分析

课程 平台 链接
Business Analytics Coursera 链接
精益数据分析 Udemy 链接
数据化决策 得到 链接

🛠️ 工具使用

Spark SQL

from pyspark.sql import SparkSession
from pyspark.sql.functions import *

# 创建 SparkSession
spark = SparkSession.builder \
    .appName("BusinessAnalytics") \
    .getOrCreate()

# 读取数据
df = spark.read.csv('sales_data.csv', header=True, inferSchema=True)

# 数据探索
df.printSchema()
df.show()

# 聚合分析
df.groupBy('category') \
    .agg(
        sum('revenue').alias('total_revenue'),
        avg('price').alias('avg_price'),
        count('order_id').alias('order_count')
    ) \
    .orderBy(desc('total_revenue')) \
    .show()

# 窗口函数
from pyspark.sql.window import Window

window_spec = Window.partitionBy('customer_id').orderBy('order_date')

df.withColumn('cumulative_revenue', 
              sum('revenue').over(window_spec)) \
  .withColumn('prev_order', 
              lag('order_date', 1).over(window_spec)) \
  .show()

# 用户留存分析
cohort_data = df.withColumn('cohort_month', 
                            trunc('first_order_date', 'MM')) \
                .withColumn('activity_month', 
                            trunc('order_date', 'MM')) \
                .withColumn('months_since_first', 
                            months_between('activity_month', 'cohort_month'))

retention = cohort_data.groupBy('cohort_month', 'months_since_first') \
    .agg(countDistinct('customer_id').alias('active_users'))

retention.show()

Pandas 性能优化

# 1. 使用 category 类型
df['category'] = df['category'].astype('category')
print(df['category'].memory_usage(deep=True))

# 2. 分块读取大数据
chunk_iter = pd.read_csv('large_file.csv', chunksize=10000)
for chunk in chunk_iter:
    process(chunk)

# 3. 使用 numba 加速
from numba import jit

@jit(nopython=True)
def fast_computation(arr):
    return arr.sum()

# 4. 并行处理
from multiprocessing import Pool

with Pool(4) as p:
    results = p.map(process_function, data_chunks)

# 5. 使用 Dask 处理超大数据
import dask.dataframe as dd

ddf = dd.read_csv('large_file.csv')
result = ddf.groupby('category')['revenue'].sum().compute()

📊 实战案例

电商销售分析

import pandas as pd

# 1. 销售趋势分析
df['order_date'] = pd.to_datetime(df['order_date'])
df['month'] = df['order_date'].dt.to_period('M')

monthly_sales = df.groupby('month')['revenue'].sum()
monthly_sales.plot(figsize=(12, 6))

# 2. 产品 ABC 分析
product_sales = df.groupby('product_id')['revenue'].sum().sort_values(ascending=False)
product_sales_cumsum = product_sales.cumsum() / product_sales.sum()

# A 类产品(贡献 80% 收入)
a_products = product_sales_cumsum[product_sales_cumsum <= 0.8].index

# 3. 客户分群(RFM)
from datetime import datetime

reference_date = df['order_date'].max()

rfm = df.groupby('customer_id').agg({
    'order_date': lambda x: (reference_date - x.max()).days,
    'order_id': 'count',
    'revenue': 'sum'
}).rename(columns={
    'order_date': 'recency',
    'order_id': 'frequency',
    'revenue': 'monetary'
})

# 分位数打分
rfm['r_score'] = pd.qcut(rfm['recency'], 4, labels=[4, 3, 2, 1])
rfm['f_score'] = pd.qcut(rfm['frequency'].rank(method='first'), 4, labels=[1, 2, 3, 4])
rfm['m_score'] = pd.qcut(rfm['monetary'], 4, labels=[1, 2, 3, 4])

rfm['rfm_score'] = rfm['r_score'].astype(int) + rfm['f_score'].astype(int) + rfm['m_score'].astype(int)

# 客户分群
def segment(rfm_score):
    if rfm_score >= 10:
        return 'Champions'
    elif rfm_score >= 7:
        return 'Loyal Customers'
    elif rfm_score >= 4:
        return 'At Risk'
    else:
        return 'Lost'

rfm['segment'] = rfm['rfm_score'].apply(segment)

用户行为漏斗

# 转化漏斗分析
funnel = pd.DataFrame({
    'stage': ['Page View', 'Add to Cart', 'Initiate Checkout', 'Purchase'],
    'users': [
        df['session_id'].nunique(),
        df[df['event'] == 'add_to_cart']['session_id'].nunique(),
        df[df['event'] == 'initiate_checkout']['session_id'].nunique(),
        df[df['event'] == 'purchase']['session_id'].nunique()
    ]
})

funnel['conversion_rate'] = funnel['users'] / funnel['users'].iloc[0]

# 可视化
import plotly.graph_objects as go

fig = go.Figure(go.Funnel(
    y = funnel['stage'],
    x = funnel['users'],
    textposition = "inside",
    textinfo = "value+percent initial"
))
fig.show()

预测建模

from sklearn.model_selection import train_test_split
from xgboost import XGBRegressor
from sklearn.metrics import mean_absolute_error, r2_score

# 特征工程
features = ['customer_age', 'days_since_last_order', 'total_orders', 
            'avg_order_value', 'days_as_customer']
X = df[features]
y = df['next_month_revenue']

# 数据分割
X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2)

# 模型训练
model = XGBRegressor(
    n_estimators=100,
    learning_rate=0.1,
    max_depth=5,
    subsample=0.8,
    colsample_bytree=0.8
)
model.fit(X_train, y_train)

# 预测与评估
y_pred = model.predict(X_test)
mae = mean_absolute_error(y_test, y_pred)
r2 = r2_score(y_test, y_pred)

print(f"MAE: {mae:.2f}")
print(f"R²: {r2:.3f}")

# 特征重要性
importance = pd.DataFrame({
    'feature': features,
    'importance': model.feature_importances_
}).sort_values('importance', ascending=False)

print(importance)

🔍 最佳实践

数据仓库设计

分层架构

graph TB
    A[ODS 操作数据层] --> B[DWD 明细数据层]
    B --> C[DWS 汇总数据层]
    C --> D[ADS 应用数据层]
    D --> E[BI 报表]
    D --> F[数据应用]

维度建模

# 星型模式示例
# 事实表:sales_fact
# 维度表:date_dim, customer_dim, product_dim, store_dim

# SQL 查询示例
"""
SELECT 
    d.year_month,
    c.customer_segment,
    p.category,
    SUM(f.revenue) as total_revenue,
    COUNT(DISTINCT f.order_id) as order_count
FROM sales_fact f
JOIN date_dim d ON f.date_id = d.date_id
JOIN customer_dim c ON f.customer_id = c.customer_id
JOIN product_dim p ON f.product_id = p.product_id
WHERE d.year = 2024
GROUP BY d.year_month, c.customer_segment, p.category
ORDER BY d.year_month, total_revenue DESC
"""

性能优化

SQL 优化

-- 1. 使用 EXPLAIN 分析
EXPLAIN SELECT * FROM orders WHERE customer_id = 100;

-- 2. 创建索引
CREATE INDEX idx_customer_order ON orders(customer_id, order_date);

-- 3. 分区表
CREATE TABLE sales_partitioned (
    order_id INT,
    order_date DATE,
    revenue DECIMAL
) PARTITION BY RANGE (YEAR(order_date)) (
    PARTITION p2020 VALUES LESS THAN (2021),
    PARTITION p2021 VALUES LESS THAN (2022),
    PARTITION p2022 VALUES LESS THAN (2023)
);

-- 4. 避免 SELECT *
SELECT order_id, customer_id, revenue FROM orders;

-- 5. 使用物化视图
CREATE MATERIALIZED VIEW mv_monthly_sales AS
SELECT 
    DATE_TRUNC('month', order_date) as month,
    SUM(revenue) as total_revenue
FROM orders
GROUP BY DATE_TRUNC('month', order_date);

🔗 更多资源

大数据

商业分析


最后更新: 2026-06-01