大数据开发学习Day4
上强度了
一、Linux/Shell 进阶
- awk 日志分析
awk 是一个强大的文本处理工具,它逐行读取文件,以字段为单位进行处理,非常适合处理结构化的文本(如 CSV、日志文件等)
awk 'pattern { action }' filename
pattern:匹配条件(行或字段满足什么条件才执行动作)
action:对匹配的行执行的操作(如打印、计算等)
- 基本打印
# 打印整行
awk '{ print $0 }' file.txt
# 打印第1和第3列
awk '{ print $1, $3 }' file.txt
# 打印行号和内容
awk '{ print NR, $0 }' file.txt
$0 整行内容
$1, $2… 第1列、第2列…(默认以空格/TAB分隔)
NR 当前行的行号(Number of Records)
- 指定分隔符
# 使用逗号作为分隔符(CSV文件)
awk -F ',' '{ print $1, $2 }' file.csv
# 使用冒号作为分隔符
awk -F ':' '{ print $1 }' /etc/passwd
# 使用正则表达式作为分隔符(空格或逗号)
awk -F '[ ,]+' '{ print $1, $2 }' file.txt
- 条件过滤
# 打印行号大于10的行
awk 'NR > 10 { print }' file.txt
# 打印第2列等于"error"的行
awk '$2 == "error" { print $0 }' log.txt
# 打印第3列大于100的行
awk '$3 > 100' file.txt
# 匹配正则表达式(包含"ERROR"的行)
awk '/ERROR/ { print }' log.txt
# 组合条件(第2列大于5且第3列小于10)
awk '$2 > 5 && $3 < 10' file.txt
- 文本处理与格式化
# 添加行号并格式化输出
awk '{ printf "%3d: %s\n", NR, $0 }' file.txt
# 修改输出分隔符
awk 'BEGIN { OFS="," } { print $1, $2, $3 }' file.txt
# 在输出中添加文本
awk '{ print "Name:", $1, "Age:", $3 }' file.txt
- 计算与统计
# 计算第3列的总和
awk '{ sum += $3 } END { print "Total:", sum }' file.txt
# 计算平均值
awk '{ sum += $3; count++ } END { print "Average:", sum/count }' file.txt
# 找出最大值
awk 'max < $3 { max = $3 } END { print "Max:", max }' file.txt
# 统计行数
awk 'END { print NR }' file.txt
- 数组用法
# 统计每个值出现的次数
awk '{ count[$1]++ } END { for (word in count) print word, count[word] }' file.txt
# 去重(打印第一次出现的行)
awk '!seen[$0]++' file.txt
任务:
看懂并执行这两个命令
自己造一个access.log测试
# 取第1列+第3列
awk '{print $1, $3}' access.log
# 第2列数字求和
awk '{sum += $2} END {print "总和:", sum}' access.log
- access.log
# 这是一个简化的日志格式:IP地址 响应时间 请求路径 状态码
192.168.1.10 25 /api/user 200
10.0.0.5 18 /api/login 200
192.168.1.20 302 /api/data 404
10.0.0.15 42 /api/order 500
192.168.1.10 36 /api/user 200
10.0.0.5 21 /api/logout 200
172.16.0.8 98 /api/upload 413
192.168.1.30 15 /api/health 200
10.0.0.20 67 /api/export 200
172.16.0.8 103 /api/import 500
# 取第1列+第3列
awk '{print $1, $3}' access.log
#结果
192.168.1.10 /api/user
10.0.0.5 /api/login
192.168.1.20 /api/data
10.0.0.15 /api/order
192.168.1.10 /api/user
10.0.0.5 /api/logout
172.16.0.8 /api/upload
192.168.1.30 /api/health
10.0.0.20 /api/export
172.16.0.8 /api/import
# 第2列数字求和
awk '{sum += $2} END {print "总和:", sum}' access.log
总和: 727
二、SQL 练习
连续类问题通用思路:排序生成行号 → 行号 - 分组内排序号 = 分组标记
1.LAG(), LEAD() 取上一行、下一行数据
2.CASE WHEN 实现行转列
3.多条件排序与分组嵌套
- 换座位(LEAD/LAG 必考)

SELECT
CASE
WHEN id % 2 = 1 AND id = (SELECT MAX(id) FROM Seat) THEN id
WHEN id % 2 = 1 THEN id + 1
ELSE id - 1
END AS id,
student
FROM Seat
ORDER BY id;
- 重新格式化部门表(行转列经典)

SELECT
id,
SUM(CASE WHEN month = 'Jan' THEN revenue END) AS Jan_Revenue,
SUM(CASE WHEN month = 'Feb' THEN revenue END) AS Feb_Revenue,
SUM(CASE WHEN month = 'Mar' THEN revenue END) AS Mar_Revenue,
SUM(CASE WHEN month = 'Apr' THEN revenue END) AS Apr_Revenue,
SUM(CASE WHEN month = 'May' THEN revenue END) AS May_Revenue,
SUM(CASE WHEN month = 'Jun' THEN revenue END) AS Jun_Revenue,
SUM(CASE WHEN month = 'Jul' THEN revenue END) AS Jul_Revenue,
SUM(CASE WHEN month = 'Aug' THEN revenue END) AS Aug_Revenue,
SUM(CASE WHEN month = 'Sep' THEN revenue END) AS Sep_Revenue,
SUM(CASE WHEN month = 'Oct' THEN revenue END) AS Oct_Revenue,
SUM(CASE WHEN month = 'Nov' THEN revenue END) AS Nov_Revenue,
SUM(CASE WHEN month = 'Dec' THEN revenue END) AS Dec_Revenue
FROM Department
GROUP BY id;
- 电影评分(综合分组 + 排序 + 窗口函数)


(SELECT name AS results
FROM MovieRating r
JOIN Users u ON r.user_id = u.user_id
GROUP BY u.user_id, name
ORDER BY COUNT(*) DESC, name
LIMIT 1)
UNION ALL
(SELECT title AS results
FROM MovieRating r
JOIN Movies m ON r.movie_id = m.movie_id
WHERE DATE_FORMAT(created_at, '%Y-%m') = '2020-02'
GROUP BY m.movie_id, title
ORDER BY AVG(rating) DESC, title
LIMIT 1);
三、今日必须掌握:
1.left /inner/right join 区别与使用场景

-
区别
- left join(左连接):返回包括左表中的所有记录和右表中联结字段相等的记录。如果右表中没有与左表匹配的记录,则右表的相关字段会用 NULL 填充。
- inner join(内连接):只返回两个表中联结字段相等的行。也就是说,只有当两个表中的联结字段值相同时,对应的记录才会出现在结果集中。
- right join(右连接):返回包括右表中的所有记录和左表中联结字段相等的记录。如果左表中没有与右表匹配的记录,则左表的相关字段会用 NULL 填充。
-
使用场景
- left join:当需要以左表为基础,获取左表的所有记录,同时获取右表中匹配的记录时使用。例如,有一个订单表和一个客户表,需要查询所有订单及其对应的客户信息,即使有些订单没有对应的客户记录(可能是数据录入问题),也希望显示订单信息,此时可以使用左连接。
- inner join:当只需要获取两个表中联结字段匹配的记录时使用。例如,在一个学生表和成绩表中,只需要查询有成绩记录的学生信息,此时可以使用内连接。
- right join:当需要以右表为基础,获取右表的所有记录,同时获取左表中匹配的记录时使用。不过在实际应用中,右连接的使用频率相对较低,因为可以通过交换表的顺序使用左连接来达到相同的效果。
-- 创建示例表
CREATE TABLE table_a (
id INT,
name VARCHAR(50)
);
CREATE TABLE table_b (
id INT,
value VARCHAR(50)
);
-- 插入示例数据
INSERT INTO table_a (id, name) VALUES (1, 'A1'), (2, 'A2'), (3, 'A3');
INSERT INTO table_b (id, value) VALUES (2, 'B2'), (3, 'B3'), (4, 'B4');
-- 使用 left join
SELECT * FROM table_a LEFT JOIN table_b ON table_a.id = table_b.id;
-- 使用 inner join
SELECT * FROM table_a INNER JOIN table_b ON table_a.id = table_b.id;
-- 使用 right join
SELECT * FROM table_a RIGHT JOIN table_b ON table_a.id = table_b.id;
2.Spark SQL 与 Hive SQL 高度一致

- 高度一致的主要表现
语法兼容性:
Spark SQL 完整支持 Hive SQL 的 DDL/DML 语法(如 CREATE TABLE、INSERT OVERWRITE)和函数(如窗口函数、UDF)。
-- Hive SQL 和 Spark SQL 均可执行
CREATE TABLE sales (id INT, product STRING);
INSERT OVERWRITE TABLE sales SELECT * FROM source_table;
元数据共享:
Spark SQL 可直接访问 Hive Metastore,实现表结构、分区等元数据的无缝共享。
// Spark 连接 Hive Metastore
val spark = SparkSession.builder()
.appName("HiveIntegration")
.config("spark.sql.warehouse.dir", "/user/hive/warehouse")
.enableHiveSupport() // 关键配置
.getOrCreate()
查询执行兼容:
支持相同的 HiveQL 查询语句,包括复杂操作(如 LATERAL VIEW、CLUSTER BY)和优化器提示(如 /+ MAPJOIN /)。
CLI工具一致性:
Spark SQL CLI 和 Beeline 提供类似 Hive Shell 的交互体验
spark-sql> SHOW TABLES; -- 命令与 Hive 相同
- 高度一致的核心原因
设计目标兼容:
Spark SQL 在设计初期就将 Hive 兼容性作为核心目标,通过 HiveContext API(Spark 2.x 前)和 Hive Support(SparkSession)实现语法解析层适配
共享元数据服务:
通过连接同一 Hive Metastore 服务(如 MySQL 或 Derby 数据库),实现表定义的统一管理
hive.metastore.uris = thrift://hive-metastore:9083
执行引擎优化:
Spark SQL 的 Catalyst 优化器兼容 Hive 的查询逻辑,同时利用 Spark 引擎的 内存计算 和 DAG调度 提升性能(比 MapReduce 快 10-100 倍)
生态整合需求:
为降低迁移成本,Spark 主动兼容 Hive 生态:
支持 Hive SerDe(如 ORC/Parquet)、兼容 Hive 安全模型(如 Ranger/Sentry)、提供 hive-site.xml 配置文件支持
- 使用场景

Hive 迁移 Spark:
无需改写 SQL,直接切换执行引擎提升性能
混合分析平台:
在同一个集群中
统一数据服务:
通过 Spark Thrift Server 提供类 Hive 的 JDBC 服务
3.cache() 懒加载机制

Spark 的 cache() 方法采用懒加载(Lazy Loading) 机制,这是 Spark 惰性计算模型的核心特性之一。其工作原理如下:
- 核心机制

延迟计算:
函数被调用时才执行计算,而非定义时
from functools import cache
@cache
def heavy_computation(n):
print(f"计算 {n}...") # 仅首次调用时执行
return n * n
print(heavy_computation(5)) # 输出"计算 5..."并返回25
print(heavy_computation(5)) # 直接返回25,无计算过程
缓存复用:
自动存储参数-结果映射(基于参数的哈希值)
- 两种实现方式
@cache
@cache
def get_data(key):
return db_query(key) # 数据库查询仅执行一次
@lru_cache(maxsize=128)
from functools import lru_cache
@lru_cache(maxsize=32)
def render_template(template_id):
# 模板渲染仅执行一次
return complex_rendering(template_id)
- 关键特性

- 使用场景
数据库查询优化
@cache
def get_user(user_id):
return User.query.get(user_id) # 避免重复数据库查询
计算缓存
@cache
def fibonacci(n):
if n < 2: return n
return fibonacci(n-1) + fibonacci(n-2) # 计算复杂度从$O(2^n)$降至$O(n)$
资源懒加载
class ImageLoader:
@cached_property # 类似机制[^2]
def high_res_image(self):
return load_from_network() # 实际使用时才加载
- 注意事项
避免可变参数
# 错误用法
@cache
def process(config: dict): ... # TypeError
# 正确用法:转为元组
@cache
def process(config_tuple: tuple): ...
缓存失效
get_user.cache_clear() # 清空所有缓存
内存泄漏风险
长期运行的进程需设置 maxsize 或定期清理
典型例子:
from functools import cache
import time
@cache
def expensive_call(param):
time.sleep(2) # 模拟耗时操作
return param.upper()
# 第一次调用(执行计算)
start = time.time()
print(expensive_call("hello")) # 输出 "HELLO" (耗时2s)
print(f"耗时: {time.time()-start:.2f}s")
# 相同参数第二次调用(使用缓存)
start = time.time()
print(expensive_call("hello")) # 立即输出 "HELLO"
print(f"耗时: {time.time()-start:.4f}s") # 耗时接近0s
4.explain() 简单看懂执行计划

在 Python 中分析数据库执行计划(如 MySQL 的 EXPLAIN),主要通过以下步骤实现:
- 获取执行计划
import pymysql
# 连接数据库
conn = pymysql.connect(host='localhost', user='root', password='123456', database='test')
cursor = conn.cursor()
# 获取执行计划
sql = """
EXPLAIN
SELECT o.*, i.*
FROM orders AS o
JOIN items AS i ON o.order_num = i.order_num
WHERE o.customer_id = 123
AND o.order_date BETWEEN '2023-01-01' AND '2023-01-02'
ORDER BY o.total_amount DESC
LIMIT 10
"""
cursor.execute(sql)
plan = cursor.fetchall() # 获取执行计划结果
# 关闭连接
cursor.close()
conn.close()
- 关键字段解读(核心关注点)

- 性能瓶颈识别
全表扫描警告
if any(row[4] == 'ALL' for row in plan):
print("⚠️ 警告:检测到全表扫描(type=ALL),考虑添加索引")
临时表警告
if any('Using temporary' in row[8] for row in plan):
print("⚠️ 警告:检测到临时表(Using temporary),优化GROUP BY/ORDER BY")
文件排序警告
if any('Using filesort' in row[8] for row in plan):
print("⚠️ 警告:检测到文件排序(Using filesort),添加ORDER BY字段索引")
- 可视化分析工具
import pandas as pd
# 转换为DataFrame分析
columns = ['id', 'select_type', 'table', 'type', 'possible_keys', 'key', 'key_len', 'ref', 'rows', 'Extra']
df = pd.DataFrame(plan, columns=columns)
# 筛选关键性能指标
print(df[['table', 'type', 'key', 'rows', 'Extra']].sort_values('rows', ascending=False))
- 优化建议
- 索引缺失
- 当 type=ALL 且 key=NULL 时,在WHERE条件字段添加索引
- 复合索引遵循最左匹配原则
- 排序优化
- 当出现 Using filesort 时,创建排序字段的索引
- 示例:ALTER TABLE orders ADD INDEX (total_amount)
- 连接优化
- 确保JOIN字段有索引(如 o.order_num, i.order_num)
- 小表驱动大表(rows小的表作为驱动表)
- 索引缺失
通过分析执行计划中的 type 和 rows 字段,可以识别出查询中需要优化的表。例如 type=ref 且 rows=1表示高效索引查找,而 type=ALL 且 rows=10000 表示全表扫描需要优化。
+----+-------------+-------+------+---------------+---------+---------+-------------------+------+-----------------------------+
| id | select_type | table | type | key | key_len | rows | Extra |
+----+-------------+-------+------+---------------+---------+---------+-----------------------------+
| 1 | SIMPLE | o | ref | customer_date | 8 | 50 | Using where; Using filesort |
| 1 | SIMPLE | i | ref | order_num | 4 | 10 | NULL |
+----+-------------+-------+------+---------------+---------+---------+-----------------------------+
优化建议:
为 o.total_amount 添加索引解决 Using filesort
检查 customer_date 索引是否包含 customer_id 和 order_date 字段
小知识点:
Spark 是懒执行,只有 show() / count() / write 这类 Action 才会真正运行
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder \
.master("local[*]") \
.appName("day4") \
.getOrCreate()
# 1. 多表JOIN(新知识点)
user_df = spark.createDataFrame(
[(1, "小明", 23), (2, "小红", 25), (3, "小李", 28)],
["id", "name", "age"]
)
score_df = spark.createDataFrame(
[(1, "math", 80), (1, "english", 85), (2, "math", 90)],
["uid", "subject", "score"]
)
# 左连接、内连接
joined = user_df.join(score_df, user_df.id == score_df.uid, how="left")
joined.show()
# 2. 创建临时视图,使用Spark SQL(新)
joined.createOrReplaceTempView("user_score")
spark.sql("""
SELECT name, age, subject, score
FROM user_score
WHERE score > 80
""").show()
# 3. 去重:dropDuplicates 指定列(新)
joined.dropDuplicates(["id"]).show()
# 4. cache 缓存(面试高频)
joined.cache()
# 5. 查看执行计划(了解即可)
joined.explain()
spark.stop()
四、算法
LeetCode 53. 最大子数组和
要求:
1.写出暴力解法
2.写出贪心优化
3.写出标准动态规划解法
理解状态转移方程:dp[i]=max(nums[i],dp[i−1]+nums[i])
def maxSubArray(nums):
dp = [0] * len(nums)
dp[0] = nums[0]
max_sum = dp[0]
for i in range(1, len(nums)):
dp[i] = max(nums[i], dp[i-1] + nums[i])
max_sum = max(max_sum, dp[i])
return max_sum
更多推荐
所有评论(0)