大数据与商业分析资源专题
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