
简介针对《大数据采集与融合技术》期末技能考核的完整解决方案文档以豆瓣2023年度书籍爬取、Flume与Kafka日志采集存储、Kettle学生成绩导入及数学排名表生成为三大主线。文档还原了考核要求与评分标准给出Python爬虫解析、Flume spooldir配置、Kafka消费者写入MySQL、Kettle Excel导入和SQL降序排名等关键实现步骤并整理常见问题的排错思路。资料为1个docx文档大小约26KB文字精炼适合具备一定编程基础、正在备考或需要课程设计参考的学生与工程师。内容按任务模块组织可对照考核项目逐项复现完整流程也能从中提炼爬虫、日志采集、ETL与数据库操作的实践方法文档还包含实验步骤与代码片段说明可帮助读者按评分要求完成实验报告和成果提交。目前已有83人学习适合期末冲刺和技能练手。1. 从期末真题到生产链路豆瓣爬虫、日志采集和成绩处理考了什么《大数据采集与融合技术》期末技能考核把豆瓣书籍爬取、日志采集与成绩处理三条链路串在一个场景里先用 Python 爬虫抓取“豆瓣 2023 年度书籍”页面的书名、作者、出版社、评分和简介再用 Flume 以 spooldir 方式监控本地日志目录把新增日志送进 Kafka最后由 Kafka 消费者写入 MySQL第三题换到另一个战场用 Kettle 把一份成绩表 Excel 导入 MySQL 并生成数学排名。三题覆盖了采集、传输、存储、加工四个环节难度不高但每一层都有生产环境会遇到的真实问题反爬识别、重复消费、偏移量提交、类型转换、排名函数的语义差异。适合手上有真题但不知从哪下手的在校生也适合想快速跑通“爬虫 Flume Kafka Kettle”最小链路的新人工程师。2. Python 爬虫解析豆瓣年度书单从 HTML 结构到 CSV 落地2.1 先弄清页面结构再写选择器豆瓣年度书单页面是服务端渲染的 HTML不是接口直出的 JSON直接请求页面本身就能拿全字段。用开发者工具查看“豆瓣 2023 年度书籍”页面源码每本书被包在一个div.item-root容器里标题在h2 a评分在.rating出版社和出版年在.publisher页数和定价在.price。这种结构下用 requests 拿全文、BeautifulSoup 做解析是最直接的组合不必上 Scrapy。当前考核只需一页数据单线程加固定间隔已经绰绰有余等以后要爬“豆瓣读书 Top 250”那种十几页的列表时再把分页循环和异常重试补上也不迟。豆瓣页面有反爬习惯高频访问会先出现 418 状态码再往后就是封禁 IP。应对方式是控制节奏而不是堆并发。把time.sleep(2)写在循环里请求间隔固定下来远比随机 UA 有用。代理池在这种低强度抓取里不是必需品代理失效本身还会引入新的不可控因素。2.2 字段映射与解析示例代码先定义字段映射表再写函数后面定位问题好对照。下表是本次考核的字段来源目标字段页面定位方式清洗处理书名div.item-root h2 a去首尾空白作者.author去掉“作者:”前缀出版社.publisher用·拆包取第一段出版年.publisher用·拆包取最后一段页数.price正则匹配(\d)页定价.price正则匹配价格数字豆瓣评分.rating取浮点数内容简介.intro压缩换行和空格完整 python 爬虫示例代码可以写成独立函数方便逐段检查import re import time import requests import pandas as pd from bs4 import BeautifulSoup HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36, Referer: https://book.douban.com/annual/2023/ } def parse_item(item: BeautifulSoup) - dict: title_node item.select_one(h2 a) author_node item.select_one(.author) pub_node item.select_one(.publisher) price_node item.select_one(.price) rating_node item.select_one(.rating) intro_node item.select_one(.intro) pub_text pub_node.get_text( , stripTrue) if pub_node else pub_parts pub_text.split(·) price_text price_node.get_text(stripTrue) if price_node else page_match re.search(r(\d)页, price_text) price_match re.search(r(\d(?:\.\d)?), price_text) rating_text rating_node.get_text(stripTrue) if rating_node else rating_match re.search(r(\d(?:\.\d)?), rating_text) return { 书名: title_node.get_text(stripTrue) if title_node else , 作者: author_node.get_text(stripTrue).replace(作者:, ) if author_node else , 出版社: pub_parts[0].strip() if len(pub_parts) 0 else , 出版年: pub_parts[-1].strip() if len(pub_parts) 1 else , 页数: page_match.group(1) if page_match else , 定价: price_match.group(1) if price_match else , 豆瓣评分: rating_match.group(1) if rating_match else , 内容简介: re.sub(r\s, , intro_node.get_text(stripTrue)) if intro_node else } def main(): url https://book.douban.com/annual/2023/ resp requests.get(url, headersHEADERS, timeout10) resp.raise_for_status() soup BeautifulSoup(resp.text, html.parser) rows [parse_item(item) for item in soup.select(div.item-root)] pd.DataFrame(rows).to_csv(douban_annual_2023.csv, indexFalse, encodingutf-8-sig) if __name__ __main__: main()这段示例代码的要点在清洗层。作者字段里的“作者:”前缀直接替换掉出版社和出版年混在同一行文本里先用·分割再取第一段和最后一段能覆盖大多数条目页数和定价都出现在.price的文本里分别用两个正则抽取比先 split 再判断更稳。最后用utf-8-sig写 CSV是为了让 Excel 直接双击打开不乱码报告截图时观感更好。2.3 分页参数与请求节奏控制年度书单单页就能拿齐数据但换到“豆瓣读书 Top 250”这类列表页时URL 会出现start游标参数每次递增 20 条。分页循环写成下面这样即可all_rows [] for offset in range(0, 251, 20): page_url fhttps://book.douban.com/top250?start{offset} resp requests.get(page_url, headersHEADERS, timeout10) if resp.status_code ! 200: time.sleep(5) continue soup BeautifulSoup(resp.text, html.parser) all_rows.extend(parse_item(item) for item in soup.select(div.item-root)) time.sleep(2.5)range(0, 251, 20)表示从第 0 条开始取到第 250 条每次跨 20 条状态码不是 200 时先 sleep 5 秒再继续不立刻退出请求成功后固定等 2.5 秒让访问频率更像人工操作。考试现场网络环境不稳定时这种“失败重试 固定间隔”的组合比任何代理配置都重要。要注意的是换页之后.item-root选择器是否仍然成立取决于目标列表页是否复用同一套模板爬之前先在浏览器里确认一次。3. Flume spooldir 采集到 Kafka配置模块拆分与启动验证3.1 为什么选 spooldir 而不是 execFlume 采集本地文件有两种常用 source。exec 走tail -F跟踪的是文件句柄进程重启或文件被轮转后会丢位置spooldir 监控整个目录把待处理文件放进去后 Flume 会逐行读取读完即将文件改名加后缀天然避免了重复读。考试题目明确指定数据源类型是 spooldir说明出题人希望看到的是“落盘文件完整后采集”的场景而不是 tail 跟随。两者的差异在答案里写清楚比只贴配置更容易拿分。Sink 端选 KafkaSink 而不是直接写文件同样有理由Kafka 扮演缓冲层采集速度和消费速度不一致时Kafka 可以暂时积压消息消费者按自己的节奏落库。所谓“缓存”不是把数据留在内存里而是让 Kafka 作为削峰填谷的中间队列。下面这份配置是完整可用的示例代码agent 名、source、channel、sink 各占一段agent.sources fileSrc agent.channels memCh agent.sinks kafkaSink agent.sources.fileSrc.type spooldir agent.sources.fileSrc.spoolDir /home/student/logs/input agent.sources.fileSrc.fileSuffix .COMPLETED agent.sources.fileSrc.deletePolicy never agent.sources.fileSrc.batchSize 100 agent.channels.memCh.type memory agent.channels.memCh.capacity 10000 agent.channels.memCh.transactionCapacity 1000 agent.sinks.kafkaSink.type org.apache.flume.sink.kafka.KafkaSink agent.sinks.kafkaSink.kafka.bootstrap.servers localhost:9092 agent.sinks.kafkaSink.kafka.topic log_topic agent.sinks.kafkaSink.flumeBatchSize 200 agent.sinks.kafkaSink.kafka.producer.acks 1 agent.sources.fileSrc.channels memCh agent.sinks.kafkaSink.channel memCh配置里最值得说的是几个容易被忽略的参数。fileSuffix决定读完后的文件后缀默认是.COMPLETED配合deletePolicy never可以保留原始文件考核现场方便复核batchSize是每个事务处理的文件行数业务日志量大时可以调到 500 以上但对本场景 100 足够。KafkaSink的flumeBatchSize控制一批发送给 Kafka 的消息条数producer.acks 1表示 leader 写入即返回吞吐和可靠性平衡。3.2 关键参数速查表把容易出错的参数单独列出来配置时一眼就能对上参数本次取值作用出错现象spoolDir/home/student/logs/input监控目录目录不存在时 agent 直接拒绝启动fileSuffix.COMPLETED读完后改名不改名可能重复采集deletePolicynever是否删除已读文件默认 never改错会丢原始文件bootstrap.serverslocalhost:9092Kafka 地址连不上时 sink 报 Connection refusedkafka.topiclog_topic目标主题主题不存在时自动创建失败则抛异常memory.capacity10000通道最大事件数超出时 source 阻塞spooldir 有一个容易踩的坑放入监控目录的文件一旦开始读就不能再原地修改追加内容。Flume 认为文件已经处理完直接改后缀新增的行不会被读到。考试时自拟日志文件应该先由脚本写完整再移动到spoolDir而不是边写边采。多个日志文件同时放入也没问题Flume 会按文件修改时间串行处理。3.3 启动顺序与端到端验证启动前先确保 Kafka 和 Zookeeper 在线主题存在。创建主题的命令kafka-topics.sh --create \ --bootstrap-server localhost:9092 \ --replication-factor 1 \ --partitions 1 \ --topic log_topic单机环境下replication-factor只能配 1分区数 1 对日志这类顺序消费场景没有压力如果要让后续消费者并行处理再增加分区数。启动 Flumeflume-ng agent \ -n agent \ -f flume-kafka.conf \ -Dflume.root.loggerINFO,console控制台刷出Event put to channel和Sink ... Event taken表示通道传输正常。验证 Flume 到 Kafka 是否通畅直接开一个终端消费者kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic log_topic \ --from-beginning能看到日志行就说明 Kafka 缓存层已经工作。注意这里用--from-beginning是为了验证历史数据如果只想看新数据可以去掉但那样容易误判为链路不通。4. Kafka 消费者落库 MySQL偏移量策略与重复消费处理4.1 消费者端为什么把 group.id 和偏移量当重点三题之中日志采集链路最容易拉开分差的是消费者落库这一步。Flume 把数据送进 Kafka 后消息不消费就不会消失消费者怎么取、取完怎么记录位置直接决定数据是否重复或丢失。Kafka 靠group.id区分不同消费组同一个组里分区只能分给一个消费者考试环境单消费者单分区group.id只要固定即可。偏移量是每个分区内消费到的位置提交偏移量之后消费者重启才会从提交点继续读。很多人在这里犯的错误是落库成功但偏移量没提交重启后重复消费全部历史消息。反过来如果先提交偏移量再写库写库失败就永久丢数据。正确顺序是先写 MySQL再提交偏移量也就是业务成功语义放在数据库上而不是放在 Kafka 的 ack 上。4.2 消费者示例代码先落库再提交偏移量下面这段示例代码用 Java 写Kafka 客户端版本为 2.x 以上都适用。建表语句放前面方便对照字段CREATE TABLE log_table ( id BIGINT AUTO_INCREMENT PRIMARY KEY, log_time VARCHAR(32), log_level VARCHAR(16), log_message TEXT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );Properties props new Properties(); props.put(bootstrap.servers, localhost:9092); props.put(group.id, log-consumer); props.put(enable.auto.commit, false); props.put(auto.offset.reset, earliest); props.put(max.poll.records, 500); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); try (KafkaConsumerString, String consumer new KafkaConsumer(props); Connection conn DriverManager.getConnection( jdbc:mysql://localhost:3306/school?useSSLfalsecharacterEncodingutf8, root, password)) { consumer.subscribe(Collections.singletonList(log_topic)); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { String value record.value(); String[] parts value.split(,, 3); if (parts.length 3) { continue; // 脏数据直接跳过不占用事务 } String sql INSERT INTO log_table(log_time, log_level, log_message) VALUES (?, ?, ?); try (PreparedStatement ps conn.prepareStatement(sql)) { ps.setString(1, parts[0].trim()); ps.setString(2, parts[1].trim()); ps.setString(3, parts[2].trim()); ps.executeUpdate(); } } consumer.commitSync(); } } catch (Exception e) { e.printStackTrace(); }这段代码的逻辑是poll 一批消息逐条解析并按逗号拆分三段写库成功后再 commitSync。enable.auto.commit设为 false所有提交位置都交给最后的 commitSync保证“先写库后提交”。split(,, 3)的第三个参数 3 表示只拆成三段日志正文里就算再有逗号也不会被误拆。parts.length 3时跳过该消息避免空指针或数组越界。4.3 消费参数对照与常见事故参数推荐值说明enable.auto.commitfalse手动提交控制丢失窗口auto.offset.resetearliest无偏移量时从头消费测试期方便max.poll.records500单次 poll 上限防处理不过来session.timeout.ms默认即可消费者心跳超时时间重复消费是 Kafka 场景下最常见的事故。即使开了手动提交commitSync 之前进程崩溃重启后仍然会重新读到这批消息。要想完全避免重复需要把 MySQL 写入操作做成幂等比如给 log_table 加一个msg_key唯一索引插入时用INSERT IGNORE或ON DUPLICATE KEY UPDATE。考试报告里能写出这一条说明是真的处理过生产问题。提示Kafka 消费者里的写入事务只是单条消息级别的。如果希望一批消息要么全部入库要么全部不提交需要用数据库事务包住整个 poll 批次并配合 Kafka 事务 API复杂度会上一个台阶考核场景不必做到那一步。5. Kettle 导入 Excel 成绩表用窗口函数生成数学排名5.1 转换链路和字段映射Kettle 这题本质是“Excel 输入 → 表输出”两段式转换。新建转换后输入步骤选“Excel 输入”输出步骤选“表输出”数据库连接指向 MySQL 的 school 库。字段映射时注意三个地方stu_no要以字符串而非数字读入否则学号前导零会丢score_math、score_english、score_chinese显式指定为 Integer 类型Excel 里最后一行“自己的学号 / 自己的姓名 / 68 / 85 / 92”要用文本格式存储避免 Kettle 把 68 识别成日期或空行跳过。映射表上把 Excel 列依次对应到name、score_math等目标列启动后看执行日志里 Finished with errors。5.2 排名 SQL 与数据完整性验证数学老师要的排名表用窗口函数写比 GROUP BY 更直接。同时计算两种排名可以发现潜在问题SELECT stu_no, name, score_math, RANK() OVER (ORDER BY score_math DESC) AS math_rank, DENSE_RANK() OVER (ORDER BY score_math DESC) AS dense_math_rank FROM score ORDER BY score_math DESC;RANK() 遇到并列成绩会跳过后续名次比如成绩为 98、98、88排名结果是 1、1、3DENSE_RANK() 则是 1、1、2。考试只要求排名表用 RANK() 还是 DENSE_RANK() 都有道理但把两种结果都写在报告里并说明并列时名次规则不同解释得越清楚越能体现对窗口函数的理解。默认排序是升序降序排名要写DESC漏掉这一项排名方向会反。数据导入没把握时用题目给出的自己的学号行验证映射。进 MySQL 执行SELECT * FROM score WHERE name 自己姓名拼音;能查到 68、85、92 三个成绩说明 Excel 列顺序和表结构对齐。确认没有重复导入就跑SELECT COUNT(*) FROM score;对比 Excel 行数。Kettle 里看到输出行数是 Excel 的两倍通常是没清理目标表导致重复插入重新执行前先 TRUNCATE 一次再跑转换排名结果才可信。本文还有配套的精品资源点击获取