简介:针对《大数据采集与融合技术》期末技能考核的完整报告,围绕豆瓣年度书籍爬取、Flume与Kafka日志采集与存储、Kettle成绩数据处理三项任务展开,是一份可直接对照练习的实践型资料。全文共 1 个 docx 文件,大小约 26KB,结构包含考核要求、实施步骤、关键代码与配置片段、问题及解决方案,便于按环节逐步完成。爬虫部分使用 Python 和 BeautifulSoup 解析豆瓣年度书单,提取书名、作者、出版社、出版年、页数、定价、评分与内容简介,并完成清洗和 CSV 存储;日志部分配置 Flume 的 spooldir 数据源监控本地目录,采集自拟日志后发往 Kafka 缓存,再通过消费者写入 MySQL;成绩处理部分使用 Kettle 将 Excel 学生成绩表导入 school 库的 score 表,并生成数学成绩降序排名表。适合具备一定编程基础、正在准备课程考核或想熟悉大数据采集全流程的学生与工程师参考,已有 83 人学习浏览,能够帮助把零散知识点串联为完整项目经验。
1. 期末考核里的三块硬功夫:爬虫、日志与数据融合
当“大数据采集与融合技术”的期末考核把豆瓣书籍爬取、日志采集和学生成绩处理放进同一套题目,多数人第一反应是“三个独立题硬拼在一起”。放到真实数据链路里看,这三块恰好覆盖了外部数据接入、系统内部运行数据收集、离线数据融合三个环节:爬虫把网页信息拉进本地,日志采集把服务运行轨迹捞起来,成绩处理则是对多来源表格做清洗和关联。三者的共同点是“采完能用”,差别在于数据源形态、体量增长速度和时效要求。
对备考读者来说,这套题真正的得分点是三件事:能否用 Python 从网页中稳定抽出结构化字段;能否说清日志从产生、采集到检索的完整链路;能否用 pandas 把多张成绩表清洗并融合成可分析结果。以下按模块展开,代码可直接运行,参数有注释解释;最后补期末验收时最容易被扣分的几个细节。
2. 豆瓣书籍爬取:从 requests 单页到 Scrapy 翻页
2.1 先选方案:requests + BeautifulSoup 还是 Scrapy
期末考核里的豆瓣读书爬取,页面量通常在 25 到 250 个条目之间,属于中轻量抓取。常见可靠的做法是requests负责请求,BeautifulSoup负责解析,最后写入 CSV。只有当老师明确要求高并发、全站采集或跨站点扩展时,才引入 Scrapy;Playwright 虽然能处理 JS 渲染,但豆瓣 Top250 这种服务端渲染页面用不上,只会白白增加资源开销。
| 方案 | 解析方式 | 并发能力 | 依赖数量 | 期末场景适合度 |
|---|---|---|---|---|
| requests + BeautifulSoup | CSS 选择器 | 串行 + 手动限速 | 2 | 高 |
| requests + lxml | XPath | 串行 + 线程池 | 2 | 高 |
| Scrapy | CSS/XPath 封装 | 异步并发 | 1 个框架 | 中(题量大时) |
| Playwright | 浏览器渲染 | 低 | 重 | 低(该页面用不上) |
如果你只是要把 Top250 拿下来,我推荐第一行;如果考核要求是“写一个分布式爬虫”,那就该走 Scrapy,因为它的调度器、去重队列天然支持后续扩展。
2.2 requests + BeautifulSoup 的最小可运行代码
import csv import time import requests from bs4 import BeautifulSoup HEADERS = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)", "Referer": "https://book.douban.com/top250", "Accept-Language": "zh-CN,zh;q=0.9", } BASE_URL = "https://book.douban.com/top250?start={}" def fetch_page(start: int) -> str: resp = requests.get(BASE_URL.format(start), headers=HEADERS, timeout=10) resp.raise_for_status() return resp.text def parse_page(html: str) -> list[dict]: soup = BeautifulSoup(html, "html.parser") books = [] # 豆瓣 Top250 的每本书在一个 tr.item 容器内 for item in soup.select("tr.item"): title_tag = item.select_one("a[href*='subject']") if not title_tag: continue rating_tag = item.select_one("span.rating_nums") books.append({ "title": title_tag.get_text(strip=True), "link": title_tag["href"], "rating": float(rating_tag.get_text(strip=True)) if rating_tag else None, }) return books def main(): result = [] for page in range(10): html = fetch_page(page * 25) result.extend(parse_page(html)) time.sleep(1) # 单位秒,控制请求频率 with open("output/douban_top250.csv", "w", encoding="utf-8-sig", newline="") as f: writer = csv.DictWriter(f, fieldnames=["title", "link", "rating"]) writer.writeheader() writer.writerows(result) print(f"共写入 {len(result)} 条记录") if __name__ == "__main__": main()逐段说明:fetch_page负责网络请求,timeout=10保证网络异常时及时报错而不是阻塞;parse_page只按容器、链接、评分三个特征取值,get_text会去掉换行和多余空格,rating_tag不存在时写入None,避免类型异常。start=0, 25, 50…是豆瓣列表的分页偏移参数,10 页正好覆盖 Top250;utf-8-sig是为了让 Excel 打开 CSV 时中文不乱码,newline=""防止写入时出现空行。
提示:
tr.item和span.rating_nums是豆瓣旧版页面的稳定特征。年份不同页面结构可能改版,跑之前先print(html[:500])看返回内容,行容器变了就优先用a[href*="subject"]做兜底。
考前最值得提醒的一件事:不要一上来就抓 10 页。先抓 1 页,校验字段和编码,再放开循环次数,这样排错成本最低。
2.3 翻页逻辑与三个必调参数
在上一段代码里,真正影响成败的是三个参数:
| 参数 | 推荐值 | 作用 |
|---|---|---|
HEADERS["User-Agent"] | 完整浏览器 UA | 空 UA 会被服务端拒接 |
time.sleep | 1~2 秒 | 限制单位时间请求次数,避免对目标站造成压力 |
timeout | 5~10 秒 | 值太小易误判,值太大卡住整个流程 |
如果题目要求做成 Scrapy,只需要保留同样的解析逻辑,换成响应对象:
import scrapy class DoubanSpider(scrapy.Spider): name = "douban250" start_urls = ["https://book.douban.com/top250?start=0"] custom_settings = { "DOWNLOAD_DELAY": 2, "FEEDS": {"output/douban.json": {"format": "json", "encoding": "utf-8"}}, } def parse(self, response): for row in response.css("tr.item"): yield { "title": row.css("a[href*='subject']::text").re_first(r"\s*(\S.*?)\s*"), "link": row.css("a[href*='subject']::attr(href)").get(), "rating": row.css("span.rating_nums::text").get(), } next_link = response.css("a.next::attr(href)").get() if next_link: yield response.follow(next_link, callback=self.parse)DOWNLOAD_DELAY相当于把time.sleep移到调度层,所有并发请求之间自动保持间隔;response.follow(next_link)会拼接相对路径并继续回调parse,这就是翻页的标准化写法。Scrapy 在期末考核中的加分点在于可以用scrapy crawl douban250 -o data/book.csv一条命令输出结果,代码里连存储逻辑都不用写。
2.4 解析字段时的三个坑
第一个坑是作者、出版社、价格这类信息常被拆在多个td里,直接用soup.select("p")会一次性拿到好几段文字。常见做法是先取整块文本,再按换行符和空格拆成author、publish、price三个字段。第二个坑是get_text(strip=True)对带空格的标题处理效果不稳定,中文标题里如果混有英文书名,最好用re.sub(r"\s+", " ", text)做一次归一化。第三个坑是评分可能缺失,写入float()前必须判断None,否则一条脏记录就会中断整批采集。这三个问题在期末答辩时是老师最爱追问的知识点。
3. 日志采集:从 ELK 选型到 Loki 轻量路径
3.1 ELK 与 Loki 两条链路的本质区别
日志采集在期末题里往往只给一句话:“采集学生成绩系统运行日志,存储并支持检索”。这句话背后其实要求你选一遍采集器、存储、展示三个角色。业界最经典的组合是 ELK(Elasticsearch + Logstash + Kibana),后来的轻量版是 Filebeat + Elasticsearch + Kibana,而热词里常被问的“ELK 是否能使用 Loki 采集日志”,本质上是问能否用 Loki 替换 Elasticsearch。
答案是:可以替换,但替换的不只是存储,检索语义和成本模型都要跟着变。
| 维度 | ELK 链路 | Loki 链路 |
|---|---|---|
| 采集器 | Filebeat / Logstash | Promtail |
| 存储索引 | 全文倒排索引 | 仅索引标签,日志原文压缩存放 |
| 查询方式 | Lucene 语法 / KQL | LogQL |
| 内存占用 | Elasticsearch 默认占用较高 | 单机最低 256MB 可跑 |
| 适合期末演示 | 可行,机器不够容易卡顿 | 更轻量,一个二进制即可 |
做课程设计时,如果电脑只有 8G 内存并同时要跑爬虫和 pandas,建议选 Loki;如果老师明确要求装 Kibana 画 Dashboard,再考虑 ELK。
3.2 用一个小脚本持续产出 JSON 行格式日志
真实系统的日志量太大,期末不好演示。更稳的做法是自己写一个日志生成器,模拟学生成绩系统的成绩查询与导入模块向/var/log/demo/追加 JSON 行式日志。
import json import random import time from datetime import datetime LOG_PATH = "/var/log/demo/score_app.log" def generate_log(line_no: int): record = { "ts": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "line": line_no, "level": random.choice(["INFO", "INFO", "WARN", "ERROR"]), "module": random.choice(["score_query", "score_import", "grade_calc"]), "cost_ms": random.randint(5, 500), "student_id": random.randint(20230001, 20230200), } return json.dumps(record, ensure_ascii=False) with open(LOG_PATH, "a", encoding="utf-8") as f: for i in range(200): f.write(generate_log(i) + "\n") f.flush() # 强制写入磁盘,保证采集器能及时读到 time.sleep(0.05)flush()是关键:文件对象默认带缓冲,不主动刷盘,Filebeat 或 Promtail 可能过好几秒才会感知到新行。line字段用来验证采集链路是否丢数据;cost_ms是为后面 Grafana 或 Kibana 画耗时分布图准备的数值字段。日志格式统一成 JSON 后,在检索端也能直接拆字段,省去解析正则。
3.3 Promtail 配置与 Loki 检索命令
Promtail 侧只需给出一个最小配置,放在/etc/promtail/config.yml:
server: http_listen_port: 9080 positions: filename: /tmp/promtail-positions.yaml clients: - url: http://localhost:3100/loki/api/v1/push scrape_configs: - job_name: score-app static_configs: - targets: [localhost] labels: job: score-app __path__: /var/log/demo/*.log__path__是 Promtail 的特殊标签,指代要监听(tail)的日志路径;job: score-app作为查询时最常用的过滤标签;positions 文件记录每次读取到文件哪个字节,程序重启后不会重复采集。启动后可执行:
# 过滤某个 job 的原始日志 curl -G http://localhost:3100/loki/api/v1/query_range \ --data-urlencode 'query={job="score-app"}' \ --data-urlencode 'limit=10' # 带关键字过滤 curl -G http://localhost:3100/loki/api/v1/query_range \ --data-urlencode 'query={job="score-app"} |= "ERROR"' \ --data-urlencode 'limit=50'LogQL 里|=表示包含匹配,!~""是正则排除,| json会把日志内容按 JSON 展开成标签再参与统计。比如想看某个模块平均耗时,可以写{job="score-app"} | json | unwrap cost_ms | avg(),这在答辩演示时会让老师眼前一亮。
3.4 Filebeat + ELK 的最小链路
如果坚持走 ELK,Filebeat 配置可以这样写:
filebeat.inputs: - type: filestream paths: - /var/log/demo/*.log parsers: - ndjson: target: "" output.elasticsearch: hosts: ["http://localhost:9200"]filestream是 Filebeat 7.13 之后的输入类型,比旧版log类型更稳定,处理日志轮转时不容易丢数据;ndjson解析器会直接把每行 JSON 解析成可检索字段,目标索引缺省为filebeat-*。期末考核的评分点通常落在“日志有没有流畅地从文件进入检索界面”,把健康检查页和查询页截图放进报告,比对组件关系口述半天更实在。
4. 学生成绩处理:pandas 多源融合与结果输出
4.1 先理清成绩数据的“乱”
学生成绩处理是典型的离线数据处理场景,难度不在算法,而在表结构不统一。常见情况是:学号字段一张表叫id、另一张叫student_id;课程名称有的写“大学英语”、有的写“大英”;缺考记录是空值而不是 0;补考成功的学生会出现两行成绩。这些数据若不清洗,直接做透视会得到错误结论。
| 原始表 | 常见叫法 | 处理后统一字段 |
|---|---|---|
| 学号 / ID | student_id, xh | student_id |
| 课程名 | 大学英语 / 大英 | course |
| 成绩 | score, 分数 | score |
这里也顺带说明 pandas 里两类核心对象:数据框(DataFrame)是二维表格,序列(Series)是一列数据,列运算本质上发生在序列上。很多错误来自把序列当数据框用,或者反过来,搞清楚两者的定义差异后再写清洗代码会更顺。
4.2 读取、清洗与合并的完整流程
import pandas as pd # 数据框的两种来源:excel 与 csv scores = pd.read_excel("data/score_input.xlsx", sheet_name="2023春季") info = pd.read_csv("data/student_info.csv", encoding="utf-8-sig") # 1) 看 schema 和缺失情况 print("scores shape:", scores.shape) print(scores.isna().sum()) # 2) 列改名统一 scores = scores.rename(columns={"学号": "student_id", "课程名称": "course"}) scores["course"] = scores["course"].str.strip() # 3) 去重:同一学生同一课程只保留最后一次记录 scores = scores.drop_duplicates(subset=["student_id", "course"], keep="last") # 4) 缺失成绩用该课程中位数填充,而不是全部填 0 med = scores.groupby("course")["score"].transform("median") scores["score"] = scores["score"].fillna(med) # 5) 课程名做映射 course_map = {"大英": "大学英语", "高数": "高等数学"} scores["course"] = scores["course"].replace(course_map) # 6) 多表关联 merged = scores.merge(info, on="student_id", how="left")关键参数含义:isna().sum()返回每列缺失值数量;drop_duplicates的keep="last"保留后写入的补考成绩;groupby(...).transform("median")返回与原表等长的序列,才能直接覆盖在缺失位置;fillna(med)用序列填充序列,pandas 会自动按索引对齐。最后一行的how="left"保证以成绩表为主表,学生信息缺失时仍能保留成绩记录。
4.3 透视表与及格率输出
统计每个班各课程平均分、每门课及格率,并输出一个多 sheet Excel。
avg_by_class = merged.pivot_table( index="class_name", columns="course", values="score", aggfunc="mean", ) def pass_rate(group): return (group["score"] >= 60).mean() pass_rate_matrix = merged.groupby(["class_name", "course"]).apply( pass_rate, include_groups=False ).unstack() # 输出一个多 sheet Excel with pd.ExcelWriter("output/score_result.xlsx", engine="openpyxl") as writer: merged.sort_values(["class_name", "course"]).to_excel( writer, sheet_name="明细", index=False) avg_by_class.round(2).to_excel(writer, sheet_name="班级均分") pass_rate_matrix.round(4).to_excel(writer, sheet_name="及格率")pivot_table的三个核心参数index、columns、values分别代表行、列和值字段;aggfunc="mean"是缺省聚合方式,改成"count"可统计人数。groupby.apply里的自定义函数逐组计算,unstack()再把结果从长表转成宽表,这种长宽转换思路在后续对接前端展示时很常用。
4.4 单机放不下的数据怎么办
期末考核很少拿上万行成绩表考你,但老师爱问“如果成绩表变成几百万行呢”。两种常见做法:一是改用pd.read_csv(..., chunksize=50000)分块读取并逐块聚合;二是直接把流程迁移到dask.dataframe,接口与 pandas 几乎一致:
import dask.dataframe as dd df = dd.read_csv("data/big_score/part-*.csv") result = df.groupby(["course"])["score"].mean() print(result.compute())dd.read_csv支持通配符,一次读入多个分片文件;compute()触发真实计算,之前只是构建计算图。这个例子可以应对“流式数据处理”和“数据处理框架”的追问,不需要去背批流一体的技术名词。要注意 dask 的groupby在compute()之前是惰性的,这一点写进报告里会更可信。
5. 期末自检:代码组织、结果文件与三个高阶技巧
5.1 一份能加分的目录结构
先把工程按三个模块分文件,让老师 5 秒能看懂:
course-design/ ├─ src/ │ ├─ douban_spider.py │ ├─ log_generator.py │ ├─ promtail-config.yml │ └─ score_process.py ├─ data/ ├─ output/ └─ README.mdREADME 里写三行运行命令:python src/douban_spider.py、python src/log_generator.py、python src/score_process.py。注意爬虫要把中间失败页的记录单独落一个fail.log,这是少数能看出工程意识的地方。
5.2 交付物自检表
| 模块 | 必交产物 | 常见失分点 |
|---|---|---|
| 豆瓣爬取 | CSV/JSON + 简单统计 | 没加限速、字段里混入换行 |
| 日志采集 | promtail/filebeat 配置 + 检索截图 | 用 print 代替文件日志,检索不到 |
| 成绩处理 | Excel 多 sheet + 结果解释 | 不检查缺失值,透视时全变 NaN |
5.3 把三个模块串起来的小技巧
期末的三段代码不需要真的互相调用,但要能说明三者关系:爬虫采集的外部书籍数据可补充成绩分析中的教材维度;成绩处理脚本本身输出运行日志,再由日志采集链路形成系统运行记录。你可以加一个总入口run_all.sh,内容只有三行带set -e的命令,这样现场演示或打包交作业都比单文件可靠。
最后在答案里附一张执行截图,爬虫结果表、日志查询结果、成绩输出 Excel 各一张,截图里带上时间,比任何文字描述都有说服力。
本文还有配套的精品资源,点击获取