大数据核心知识笔记
第一部分大数据生态系统概览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( 实时推荐系统就绪)

相关新闻

如何3步实现智能图片分层:Layerdivider的终极效率指南

如何3步实现智能图片分层:Layerdivider的终极效率指南

如何3步实现智能图片分层:Layerdivider的终极效率指南 【免费下载链接】layerdivider A tool to divide a single illustration into a layered structure. 项目地址: https://gitcode.com/gh_mirrors/la/layerdivider 在当今数字设计领域,你是否…

2026/8/1 1:36:23 阅读更多 →
TSB空间斩特效:本地部署与游戏开发集成实践

TSB空间斩特效:本地部署与游戏开发集成实践

这次我们来看一个 TSB 自定义技能项目,重点是在本地环境中实现"空间斩"特效的代码级部署和功能验证。这个项目已经开源,代码可直接获取,适合想要在游戏开发、特效制作或动画生成中集成自定义技能效果的开发者。从项目标题看&#x…

2026/8/1 1:36:23 阅读更多 →
FreeRTOS链表实现与优化解析

FreeRTOS链表实现与优化解析

1. FreeRTOS链表实现深度解析在嵌入式实时操作系统领域,链表是最基础也最重要的数据结构之一。作为FreeRTOS的核心组件,其链表实现方式直接影响着任务调度、内存管理和IPC机制的效率。我第一次在STM32F103上移植FreeRTOS时,就曾因为对链表理解…

2026/8/1 1:35:23 阅读更多 →

最新新闻

易语言AI代码生成工具:中文编程的智能辅助实践指南

易语言AI代码生成工具:中文编程的智能辅助实践指南

这次我们来看一个很有意思的项目——易语言AI写代码。对于很多习惯使用易语言的开发者来说,这个工具可能是个福音。易语言作为一门中文编程语言,在国内有相当广泛的用户基础,但长期以来缺乏现代化的AI辅助编程支持。现在有了AI写代码的能力&a…

2026/8/1 2:12:37 阅读更多 →
TI杯电赛三轮两驱智能车底盘安装与调试全攻略

TI杯电赛三轮两驱智能车底盘安装与调试全攻略

这次我们来详细拆解TI杯电赛智能车的三轮两驱底盘安装全过程。对于参加电赛的同学来说,底盘安装是智能车制作的第一步,也是最关键的基础环节。三轮两驱结构因其简单可靠、成本低廉,成为众多电赛队伍的首选方案。这种底盘结构最大的特点是两个…

2026/8/1 2:12:37 阅读更多 →
SolidWorks_动画模拟与仿真15_仿真结果可视化

SolidWorks_动画模拟与仿真15_仿真结果可视化

仿真结果可视化:图解、矢量图与动态高亮的工程实践本文深入探讨工程仿真(FEA/CFD)结果可视化的核心技术,从数据映射到动态交互,完整呈现应力、位移等物理场的图形化表达方案,并附赠可直接运行的Python代码示…

2026/8/1 2:12:37 阅读更多 →
Transformer大模型实战:从自注意力原理到多模态应用开发

Transformer大模型实战:从自注意力原理到多模态应用开发

在深度学习领域,Transformer架构彻底改变了自然语言处理乃至多模态任务的格局。无论是BERT、GPT系列还是最新的多模态大模型,其核心都离不开Transformer。然而,许多开发者在理论学习与工程落地之间仍存在断层:原理看似复杂&#x…

2026/8/1 2:12:37 阅读更多 →
FilePulse:Kafka Connect 的“智能文件网关”,重塑数据接入新范式

FilePulse:Kafka Connect 的“智能文件网关”,重塑数据接入新范式

在构建实时数据湖的过程中,文件接入往往是第一道难关。传统的 FileStreamSource 仅能做简单的“文本搬运”,面对复杂格式往往捉襟见肘。本文将深入介绍 FilePulse,一款功能强大的 Kafka Connect 源连接器。它不仅能实时监控目录变化&#xff…

2026/8/1 2:12:37 阅读更多 →
Unity与Socket构建实时远程监控系统:从原理到实践

Unity与Socket构建实时远程监控系统:从原理到实践

1. 项目概述:为什么选择Unity与Socket构建远程监控? 如果你正在寻找一个既能深入理解网络通信,又能产出可视化、可交互成果的实战项目,那么用Unity结合Socket搭建一个实时画面传输的远程监控系统,绝对是一个绝佳的选择…

2026/8/1 2:11:36 阅读更多 →

日新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/1 0:00:48 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/1 0:00:48 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/1 0:00:48 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/31 1:03:03 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/29 14:34:28 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/31 4:19:39 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/1 0:00:48 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/1 0:00:48 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/1 0:00:48 阅读更多 →