大数据核心知识笔记

大数据核心知识笔记
第一部分大数据生态系统概览1. 什么是大数据文字解释大数据是指无法用传统数据库工具处理的海量数据集合。它有著名的5V特征Volume大量数据量巨大TB → PB → EBVelocity高速产生和变化速度快实时流数据Variety多样数据类型多样结构化、半结构化、非结构化Value价值价值密度低但总量价值高Veracity真实性数据质量参差不齐代码描述用Python模拟生成1GB数据展示Volumepythonimport random import csv # 生成100万条用户行为数据约200MB def generate_big_data(filenameuser_behavior.csv, rows1000000): 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 f2024-{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经典案例pythonfrom 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(processes4) 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(), keylambda x: -x[1]): print(f {word}: {count})3. Hadoop vs Spark文字解释特性Hadoop MapReduceApache Spark数据处理方式磁盘读写慢内存计算快100倍编程语言Java为主Java/Scala/Python/R适用场景批量离线处理实时流处理批处理容错机制重新计算整个任务RDD血缘关系精确恢复代码描述Spark实现单词统计对比上面Hadoop的代码pythonfrom pyspark import SparkContext, SparkConf # 创建Spark上下文 conf SparkConf().setAppName(WordCount).setMaster(local[*]) sc SparkContext(confconf) # 读取数据可以是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. HDFSHadoop分布式文件系统文字解释HDFS是专为大文件设计GB/TB级别的分布式文件系统核心设计数据分块默认128MB/块大文件切割存储副本机制默认3副本保证容错主从架构NameNode元数据 DataNode实际数据一次写入多次读取不支持文件修改适合批处理代码描述用Python模拟HDFS的读写过程pythonimport hashlib import random class HDFS_Simulator: def __init__(self, block_size128, replication3): 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 fnode_{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 ✅ 块 {i1}: {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 ✅ 读取块 {i1} 从 {node_id}) break return all_data.decode(utf-8, errorsignore) # 测试 hdfs HDFS_Simulator(block_size1) # 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窗口计算pythonimport 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(maxlen1000) # 最多保留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_size5): # 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_size3) # 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(targetproducer, daemonTrue).start() Thread(targetconsumer, daemonTrue).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 (dt2024-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引擎文字解释区别于HivePresto是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训练线性回归模型pythonfrom 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, headerTrue, inferSchemaTrue) # 特征工程将多列合并为特征向量 feature_cols [ad_spend, website_visits, social_media_budget] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) data assembler.transform(data) # 划分训练集和测试集 train, test data.randomSplit([0.8, 0.2], seed42) # 训练线性回归模型 lr LinearRegression(featuresColfeatures, labelColsales) model lr.fit(train) # 预测并评估 predictions model.transform(test) evaluator RegressionEvaluator(labelColsales, metricNamermse) 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文字解释追踪数据从源头到消费的整个生命周期回答数据从哪里来经过哪些处理被谁使用代码描述简单数据血缘追踪系统pythonclass 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_size10) 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( 实时推荐系统就绪)