大数据核心知识笔记

📅 2026/8/1 1:40:07 👁️ 阅读次数 📝 编程学习
大数据核心知识笔记

📚 第一部分:大数据生态系统概览

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 MapReduceApache 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/TrinoMPP架构,查询快
数据存储(海量冷数据)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("🎉 实时推荐系统就绪!")