📚 第一部分:大数据生态系统概览
1. 什么是大数据?
文字解释:
大数据是指无法用传统数据库工具处理的海量数据集合。它有著名的"5V"特征:
Volume(大量):数据量巨大(TB → PB → EB)
Velocity(高速):产生和变化速度快(实时流数据)
Variety(多样):数据类型多样(结构化、半结构化、非结构化)
Value(价值):价值密度低,但总量价值高
Veracity(真实性):数据质量参差不齐
代码描述:用Python模拟生成1GB数据(展示Volume)
python
import random import csv # 生成100万条用户行为数据(约200MB) def generate_big_data(filename="user_behavior.csv", rows=1000000): with open(filename, 'w', newline='') as f: writer = csv.writer(f) writer.writerow(['user_id', 'timestamp', 'action', 'product_id', 'price']) for i in range(rows): writer.writerow([ random.randint(1, 100000), # user_id f"2024-{random.randint(1,12):02d}-{random.randint(1,28):02d}", random.choice(['click', 'view', 'purchase', 'add_to_cart']), random.randint(1, 10000), round(random.uniform(1.99, 999.99), 2) ]) print(f"✅ 生成了 {rows} 条数据,文件大小: {os.path.getsize(filename)/1024/1024:.2f} MB") generate_big_data() # 运行试试!你的电脑会卡吗?🔥 第二部分:分布式计算框架
2. MapReduce 编程模型
文字解释:
MapReduce是Google提出的分布式计算模型,核心思想是"分而治之":
Map阶段:将数据拆分成多个小块,并行处理,生成键值对
Shuffle阶段:自动将相同Key的数据聚合到一起
Reduce阶段:对聚合后的数据进行汇总计算
就像把一堆拼图(Map)分给100个人同时拼,然后再把拼好的部分组合起来(Reduce)!
代码描述:用Python实现单词计数(MapReduce经典案例)
python
from collections import defaultdict import multiprocessing as mp # Map函数:将文本拆分成单词 def map_function(text_chunk): word_count = defaultdict(int) for word in text_chunk.split(): word = word.lower().strip('.,!?') if word: word_count[word] += 1 return dict(word_count) # Reduce函数:合并统计结果 def reduce_function(mapped_results): final_count = defaultdict(int) for result in mapped_results: for word, count in result.items(): final_count[word] += count return dict(final_count) # 模拟分布式处理 def mapreduce_demo(texts): # Map阶段:并行处理 with mp.Pool(processes=4) as pool: mapped = pool.map(map_function, texts) # Reduce阶段:合并结果 result = reduce_function(mapped) return result # 测试 texts = [ "Hello world Hello Hadoop", "MapReduce is powerful MapReduce", "Hello again world" ] result = mapreduce_demo(texts) print("📊 单词统计结果:") for word, count in sorted(result.items(), key=lambda x: -x[1]): print(f" {word}: {count}")3. Hadoop vs Spark
文字解释:
| 特性 | Hadoop MapReduce | Apache Spark |
|---|---|---|
| 数据处理方式 | 磁盘读写(慢) | 内存计算(快100倍) |
| 编程语言 | Java为主 | Java/Scala/Python/R |
| 适用场景 | 批量离线处理 | 实时流处理+批处理 |
| 容错机制 | 重新计算整个任务 | RDD血缘关系(精确恢复) |
代码描述:Spark实现单词统计(对比上面Hadoop的代码)
python
from pyspark import SparkContext, SparkConf # 创建Spark上下文 conf = SparkConf().setAppName("WordCount").setMaster("local[*]") sc = SparkContext(conf=conf) # 读取数据(可以是HDFS、本地文件等) text_file = sc.textFile("data.txt") # 一行代码完成单词计数! word_counts = text_file.flatMap(lambda line: line.split()) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a + b) # 收集结果 for word, count in word_counts.collect(): print(f"{word}: {count}") sc.stop()💡 关键差异:
Spark的
reduceByKey会自动在本地先做一次聚合(Map端聚合),减少网络传输RDD(弹性分布式数据集)支持懒加载,只在Action操作时真正计算
💾 第三部分:分布式存储系统
4. HDFS(Hadoop分布式文件系统)
文字解释:
HDFS是专为大文件设计(GB/TB级别)的分布式文件系统,核心设计:
数据分块:默认128MB/块,大文件切割存储
副本机制:默认3副本,保证容错
主从架构:NameNode(元数据) + DataNode(实际数据)
一次写入,多次读取:不支持文件修改(适合批处理)
代码描述:用Python模拟HDFS的读写过程
python
import hashlib import random class HDFS_Simulator: def __init__(self, block_size=128, replication=3): self.block_size = block_size * 1024 * 1024 # 转换为字节 self.replication = replication self.name_node = {} # 文件名 -> [块列表] self.data_nodes = {} # 节点ID -> {块ID: 数据} self.node_count = 5 def write_file(self, filename, data): """模拟文件写入""" data_bytes = data.encode('utf-8') total_size = len(data_bytes) block_count = (total_size + self.block_size - 1) // self.block_size print(f"📝 写入文件: {filename}") print(f" 总大小: {total_size/1024/1024:.2f} MB") print(f" 分块数: {block_count}") block_list = [] for i in range(block_count): # 切分数据块 start = i * self.block_size end = min(start + self.block_size, total_size) block_data = data_bytes[start:end] # 生成块ID block_id = hashlib.md5(f"{filename}_{i}".encode()).hexdigest()[:8] # 存储副本(模拟3副本) for j in range(self.replication): node_id = f"node_{random.randint(1, self.node_count)}" if node_id not in self.data_nodes: self.data_nodes[node_id] = {} self.data_nodes[node_id][block_id] = block_data block_list.append(block_id) print(f" ✅ 块 {i+1}: {block_id} (大小: {len(block_data)/1024:.2f} KB)") self.name_node[filename] = block_list print(f"✅ 文件写入完成!") def read_file(self, filename): """模拟文件读取""" if filename not in self.name_node: print(f"❌ 文件 {filename} 不存在") return None block_list = self.name_node[filename] print(f"📖 读取文件: {filename}") print(f" 块数: {len(block_list)}") all_data = b'' for i, block_id in enumerate(block_list): # 从任意DataNode读取(模拟负载均衡) for node_id, blocks in self.data_nodes.items(): if block_id in blocks: data = blocks[block_id] all_data += data print(f" ✅ 读取块 {i+1} 从 {node_id}") break return all_data.decode('utf-8', errors='ignore') # 测试 hdfs = HDFS_Simulator(block_size=1) # 1MB块大小用于测试 data = "Hello HDFS! " * 100000 # 约1.8MB数据 hdfs.write_file("test.txt", data) result = hdfs.read_file("test.txt") print(f"\n📄 读取内容前100字符: {result[:100]}...")5. HBase(列式存储数据库)
文字解释:
HBase是基于HDFS的NoSQL列式数据库,适合随机读写大表:
行键(Row Key):唯一标识,按字典序排序
列族(Column Family):逻辑分组,需要预定义
单元格(Cell):存储具体值,带时间戳版本
特点:支持上亿行 × 百万列的稀疏表
代码描述:HBase Shell操作示例
bash
# HBase Shell命令 hbase shell # 创建表(users表,有info和behavior两个列族) create 'users', 'info', 'behavior' # 插入数据 put 'users', 'user_1001', 'info:name', 'Alice' put 'users', 'user_1001', 'info:age', '28' put 'users', 'user_1001', 'behavior:last_login', '2024-01-15' # 批量查询(Scan) scan 'users', {STARTROW => 'user_1000', LIMIT => 10} # 单行查询(Get) get 'users', 'user_1001' # 删除列 delete 'users', 'user_1001', 'info:age'⚡ 第四部分:流式计算与实时处理
6. Kafka + Flink 实时处理
文字解释:
Kafka:分布式消息队列,像"数据管道",支持高吞吐量的发布订阅
Flink:真正的流式计算引擎(vs Spark Streaming的微批次),毫秒级延迟
经典架构:
text
数据源 → Kafka(消息队列) → Flink(实时计算) → 数据库/可视化
代码描述:用Python模拟Kafka生产和消费 + Flink窗口计算
python
import time import random from collections import deque from threading import Thread # 模拟Kafka class KafkaTopic: def __init__(self, topic_name): self.topic_name = topic_name self.messages = deque(maxlen=1000) # 最多保留1000条 def produce(self, message): self.messages.append(message) print(f"📤 [{self.topic_name}] 生产: {message}") def consume(self): if self.messages: return self.messages.popleft() return None # 模拟Flink流处理 class FlinkStreamProcessor: def __init__(self, window_size=5): # 5秒窗口 self.window_size = window_size self.window_data = [] self.last_window_time = time.time() def process(self, message): current_time = time.time() self.window_data.append(message) # 每5秒触发一次窗口计算 if current_time - self.last_window_time >= self.window_size: self.compute_window() self.window_data = [] self.last_window_time = current_time def compute_window(self): if not self.window_data: return # 假设数据是 user_id, action, amount # 计算窗口内的统计信息 total_amount = sum(item['amount'] for item in self.window_data) action_count = {} for item in self.window_data: action_count[item['action']] = action_count.get(item['action'], 0) + 1 print(f"\n📊 [窗口统计] 共 {len(self.window_data)} 条数据") print(f" 总金额: ${total_amount:.2f}") print(f" 行为分布: {action_count}") print("-" * 40) # 模拟数据流 def simulate_data_stream(): topic = KafkaTopic("user_actions") processor = FlinkStreamProcessor(window_size=3) # 3秒窗口 # 启动生产者线程 def producer(): actions = ['click', 'purchase', 'view', 'add_to_cart'] while True: message = { 'user_id': random.randint(1, 100), 'action': random.choice(actions), 'amount': round(random.uniform(1, 100), 2) if random.random() > 0.7 else 0 } topic.produce(message) time.sleep(random.uniform(0.2, 0.8)) # 启动消费者线程(Flink) def consumer(): while True: message = topic.consume() if message: processor.process(message) time.sleep(0.1) # 启动线程 Thread(target=producer, daemon=True).start() Thread(target=consumer, daemon=True).start() # 运行15秒 time.sleep(15) print("\n✅ 流处理模拟结束") # 运行模拟 simulate_data_stream()🗄️ 第五部分:数据仓库与查询引擎
7. Hive(数据仓库工具)
文字解释:
Hive将SQL语句转换为MapReduce/Spark作业,让数据分析师可以用SQL处理大数据:
元数据存储:表结构、分区信息(存储在MySQL中)
数据存储:实际数据在HDFS上
支持分区:提高查询效率(如按日期分区)
代码描述:Hive建表和查询示例
sql
-- 创建Hive表(外部表,数据在HDFS) CREATE EXTERNAL TABLE user_logs ( user_id INT, action STRING, product_id INT, price DOUBLE ) PARTITIONED BY (dt STRING) -- 按日期分区 ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' STORED AS TEXTFILE LOCATION '/data/user_logs'; -- 加载数据(从HDFS移动文件到表目录) LOAD DATA INPATH '/raw_data/2024-01-15.log' INTO TABLE user_logs PARTITION (dt='2024-01-15'); -- 数据分析查询(转为MapReduce作业) SELECT action, COUNT(*) AS cnt, AVG(price) AS avg_price FROM user_logs WHERE dt = '2024-01-15' AND price > 0 GROUP BY action ORDER BY cnt DESC; -- 创建分区表优化查询 CREATE TABLE user_behavior_partitioned ( user_id INT, behavior STRING ) PARTITIONED BY (dt STRING) STORED AS PARQUET; -- 列式存储,压缩率高
8. Presto/Trino(分布式SQL引擎)
文字解释:
区别于Hive:Presto是MPP(大规模并行处理)引擎,不依赖HDFS存储
特点:支持联邦查询(同时查询Hive、MySQL、Kafka等)
速度:比Hive快5-10倍(适合交互式查询)
代码描述:Presto查询示例
sql
-- 跨数据源联合查询 SELECT u.user_name, o.order_id, o.amount, o.order_time FROM hive.default.users u JOIN mysql.default.orders o ON u.user_id = o.user_id WHERE o.order_time > DATE '2024-01-01' AND u.country = 'China'; -- 实时查询(连接Kafka) SELECT user_id, COUNT(*) AS click_count FROM kafka.default.click_stream WHERE _timestamp > CURRENT_TIMESTAMP - INTERVAL '5' MINUTE GROUP BY user_id HAVING COUNT(*) > 10; -- 高活跃用户
🤖 第六部分:机器学习与大数据结合
9. MLlib + 分布式训练
文字解释:
大数据平台上的机器学习库,支持:
分布式算法:线性回归、随机森林、K-Means等
特征工程:标准化、PCA、TF-IDF等
模型部署:支持导出为PM模型
代码描述:Spark MLlib训练线性回归模型
python
from pyspark.sql import SparkSession from pyspark.ml.regression import LinearRegression from pyspark.ml.feature import VectorAssembler from pyspark.ml.evaluation import RegressionEvaluator # 创建Spark会话 spark = SparkSession.builder.appName("MLDemo").getOrCreate() # 准备数据(假设有10万条数据) data = spark.read.csv("sales_data.csv", header=True, inferSchema=True) # 特征工程:将多列合并为特征向量 feature_cols = ['ad_spend', 'website_visits', 'social_media_budget'] assembler = VectorAssembler(inputCols=feature_cols, outputCol="features") data = assembler.transform(data) # 划分训练集和测试集 train, test = data.randomSplit([0.8, 0.2], seed=42) # 训练线性回归模型 lr = LinearRegression(featuresCol="features", labelCol="sales") model = lr.fit(train) # 预测并评估 predictions = model.transform(test) evaluator = RegressionEvaluator(labelCol="sales", metricName="rmse") rmse = evaluator.evaluate(predictions) print(f"✅ 模型训练完成") print(f"📊 权重系数: {model.coefficients}") print(f"📊 截距: {model.intercept}") print(f"📊 RMSE: {rmse:.2f}") # 批量预测新数据 new_data = spark.createDataFrame([ (1000, 5000, 200), (2000, 8000, 300) ], feature_cols) new_data = assembler.transform(new_data) predictions = model.transform(new_data) predictions.show()🛠️ 第七部分:数据治理与监控
10. 数据血缘(Data Lineage)
文字解释:
追踪数据从源头到消费的整个生命周期,回答"数据从哪里来?经过哪些处理?被谁使用?"
代码描述:简单数据血缘追踪系统
python
class DataLineage: def __init__(self): self.graph = {} # 节点关系图 def add_transformation(self, source, target, operation): """记录数据转换关系""" if source not in self.graph: self.graph[source] = [] self.graph[source].append({ 'target': target, 'operation': operation, 'timestamp': time.time() }) print(f"🔗 记录血缘: {source} --{operation}--> {target}") def trace_source(self, target): """追溯数据源头""" print(f"\n🔍 追溯 {target} 的数据来源:") current = target path = [current] def find_parent(node): for parent, children in self.graph.items(): for child in children: if child['target'] == node: path.append(parent) find_parent(parent) return find_parent(current) path.reverse() for i, node in enumerate(path): print(f" {' ' * i}└── {node}") # 使用示例 lineage = DataLineage() lineage.add_transformation("user_logs", "cleaned_logs", "filter_null") lineage.add_transformation("cleaned_logs", "user_behavior", "aggregate") lineage.add_transformation("user_behavior", "sales_report", "join_with_orders") lineage.trace_source("sales_report")📊 性能对比与最佳实践
| 场景 | 推荐工具 | 理由 |
|---|---|---|
| 离线批处理(TB级) | Spark/Hadoop | 稳定可靠,成本低 |
| 实时流处理(毫秒级) | Flink | 真正实时,高吞吐 |
| 交互式查询(秒级响应) | Presto/Trino | MPP架构,查询快 |
| 数据存储(海量冷数据) | HDFS + Parquet | 压缩率高,成本低 |
| 随机读写(实时更新) | HBase | 支持上亿行随机读写 |
| 消息队列(解耦系统) | Kafka | 高吞吐,持久化 |
🎯 综合实践:构建实时电商推荐系统
把上面所有知识串起来,实现一个简单的推荐系统!
python
""" 电商实时推荐系统架构: 1. Kafka接收用户点击流 2. Flink实时计算用户画像 3. Spark离线训练推荐模型 4. Redis存储实时特征 5. HBase存储用户历史 """ # 这里只展示核心流程(伪代码) class RealtimeRecommendationSystem: def __init__(self): self.kafka = KafkaTopic("user_click") self.flink = FlinkStreamProcessor(window_size=10) self.redis = {} # 模拟Redis缓存 self.model = None # 预训练模型 def process_user_click(self, user_id, product_id): # 实时更新用户特征 self.update_user_profile(user_id, product_id) # 实时推荐 recommendations = self.get_recommendations(user_id) return recommendations def update_user_profile(self, user_id, product_id): # Flink实时计算 self.flink.process({ 'user_id': user_id, 'product_id': product_id, 'action': 'click', 'timestamp': time.time() }) def get_recommendations(self, user_id): # 1. 从Redis获取用户实时特征 user_features = self.redis.get(user_id, {}) # 2. 从HBase获取用户历史 history = self.hbase.get(user_id, 'history') # 3. 使用Spark ML模型预测 candidates = self.model.predict(user_features, history) return candidates[:10] # 返回Top10推荐 print("🎉 实时推荐系统就绪!")