大数据开发学习Day9
·
一、Linux / Shell 新知识点
今日新:cut + 字段截取 + 日志实战
任务:
有如下格式日志 login.log
2025-04-09 10:00:00,uid100,login_success
2025-04-09 10:01:00,uid101,login_fail
2025-04-09 10:02:00,uid100,login_success
要求一行命令:
按逗号分隔
提取第 2 列(uid)
去重统计出现次数
按次数倒序
cut -d ',' -f 2 login.log | sort | uniq -c | sort -nr
cut -d ‘,’ -f 2:按逗号分隔,取第 2 列
适合处理 CSV / 日志类结构化数据
与 awk 相比更轻量、更快
二、SQL
1.用户连续登录天数统计

WITH t1 AS (
SELECT
user_id,
login_date,
-- 按用户分组,日期排序,生成行号
ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY login_date) AS rn
FROM user_login
),
t2 AS (
SELECT
user_id,
login_date,
-- 连续日期的分组标记:日期 - 行号 = 相同值
DATE_SUB(login_date, INTERVAL rn DAY) AS grp
FROM t1
),
t3 AS (
SELECT
user_id,
MIN(login_date) AS start_date,
MAX(login_date) AS end_date,
COUNT(*) AS continuous_days
FROM t2
GROUP BY user_id, grp
)
SELECT *
FROM t3
WHERE continuous_days >= 3;
连续日期通用解法:行号差值分组法
DATE_SUB(login_date, INTERVAL rn DAY)
连续日期的结果相同,形成分组
先分组求起止日期与天数,再筛选≥3 天
面试高频:连续登录、连续签到、连续消费
2.各部门薪资占比

WITH dept_total AS (
SELECT
department_id,
SUM(salary) AS total_salary
FROM employee
GROUP BY department_id
)
SELECT
e.id,
e.name,
e.department_id,
e.salary,
ROUND(e.salary / d.total_salary * 100, 2) AS salary_percent
FROM employee e
JOIN dept_total d ON e.department_id = d.department_id;
先算部门总和,再关联个人
占比公式:个人薪资 / 部门总薪资 ×100
ROUND(…,2) 保留两位小数
典型场景:占比、份额、贡献度统计
3.每用户最大连续交易天数

WITH t1 AS (
SELECT
user_id,
trans_date,
ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY trans_date) AS rn
FROM trans
),
t2 AS (
SELECT
user_id,
trans_date,
DATE_SUB(trans_date, INTERVAL rn DAY) AS grp
FROM t1
),
t3 AS (
SELECT
user_id,
COUNT(*) AS days
FROM t2
GROUP BY user_id, grp
)
SELECT
user_id,
MAX(days) AS max_continuous_days
FROM t3
GROUP BY user_id;
与连续登录同一套模板
先分组求每组连续天数
再用 MAX() 取每个用户的最大值
是数仓、用户增长最常用统计逻辑
三、PySpark 新内容
广播 join + 数据倾斜实战 + 多重聚合 + 临时视图函数
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder \
.master("local[*]") \
.appName("day9") \
.getOrCreate()
# 大表
big_df = spark.createDataFrame([
(1, "click"), (2, "view"), (1, "buy"), (3, "click")
], ["uid", "event"])
# 小表
small_df = spark.createDataFrame([
(1, "小明"), (2, "小红"), (3, "小李")
], ["uid", "name"])
# ==================== 新知识点1:广播Join(解决数据倾斜神器)====================
# 自动广播小表
df_join = big_df.join(F.broadcast(small_df), on="uid", how="left")
df_join.show()
# ==================== 新知识点2:多层级聚合 ====================
df_join.groupBy("name", "event") \
.agg(F.count("*").alias("cnt")) \
.groupBy("name") \
.agg(F.sum("cnt").alias("total")) \
.show()
# ==================== 新知识点3:查看执行计划 ====================
df_join.explain(True)
spark.stop()
运行结果:
+---+-----+----+
|uid|event|name|
+---+-----+----+
| 1|click|小明|
| 2| view|小红|
| 1| buy|小明|
| 3|click|小李|
+---+-----+----+
+----+-----+
|name|total|
+----+-----+
|小明| 2|
|小红| 1|
|小李| 1|
+----+-----+
1.广播 join 适用场景(大表 join 小表)
小表发送到所有节点,大表本地关联
无 shuffle,极大提升速度
适用:大表 Join 小表(数据倾斜神器)
2.广播能显著减少 shuffle、缓解倾斜
3.多层级聚合写法
先按 name+event 统计
再按 name 汇总
4.explain 看懂简单执行计划
查看逻辑计划、物理计划
面试常问:如何判断是否广播、是否 shuffle
四、算法
LeetCode 239. 滑动窗口最大值
理解单调队列思想
能写出 O (n) 解法
知道暴力解法为什么超时
from collections import deque
def maxSlidingWindow(nums, k):
q = deque()
res = []
for i, num in enumerate(nums):
# 左边超出窗口,移除
while q and q[0] <= i - k:
q.popleft()
# 保持队列递减,小的直接弹出
while q and nums[q[-1]] <= num:
q.pop()
q.append(i)
# 窗口形成后开始记录
if i >= k - 1:
res.append(nums[q[0]])
return res
更多推荐
所有评论(0)