Storm Crawler爬取结果存储位置及数据库存储方案咨询
Storm Crawler 爬取结果存储与数据库适配指南
一、默认爬取结果存储位置
Storm Crawler 未配置 Solr/ES 时,爬取的原始数据(网页内容、元数据等)默认暂存在本地文件系统临时目录或内存队列中,这类存储是临时的,重启拓扑后数据会丢失,不适合长期保存。默认临时路径通常在项目的./target目录下,但这不是稳定的持久化方案。
二、配置数据库存储的步骤
作为新手,你可以通过自定义 Storm Bolt 组件实现数据库存储——Storm Crawler 基于 Storm 拓扑架构,核心由 Spout(数据输入)和 Bolt(数据处理)组成,我们只需添加自定义 Bolt 来接收解析后的文档,写入数据库即可。
1. 选择适配的数据库
结合后续对接搜索工具的需求,推荐优先选择:
- 关系型数据库:PostgreSQL/MySQL,适合结构化数据存储,方便后续导出转换
- 文档型数据库:MongoDB,适配非结构化网页内容,导入 Typesense/Flex Search 更灵活
2. 编写自定义存储 Bolt(Java 示例)
Storm Crawler 主要用 Java 开发,以下是对接 MySQL 的自定义 Bolt 代码:
import org.apache.storm.topology.BasicOutputCollector; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseBasicBolt; import org.apache.storm.tuple.Tuple; import com.digitalpebble.stormcrawler.Metadata; import com.digitalpebble.stormcrawler.util.ConfUtils; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.SQLException; public class DatabaseStorageBolt extends BaseBasicBolt { private Connection dbConn; private String dbUrl; private String dbUser; private String dbPwd; @Override public void prepare(java.util.Map conf, org.apache.storm.task.TopologyContext context) { // 从配置文件读取数据库参数 dbUrl = ConfUtils.getString(conf, "db.url"); dbUser = ConfUtils.getString(conf, "db.user"); dbPwd = ConfUtils.getString(conf, "db.password"); try { Class.forName("com.mysql.cj.jdbc.Driver"); dbConn = DriverManager.getConnection(dbUrl, dbUser, dbPwd); } catch (ClassNotFoundException | SQLException e) { throw new RuntimeException("数据库连接失败", e); } } @Override public void execute(Tuple tuple, BasicOutputCollector collector) { String url = tuple.getStringByField("url"); String content = tuple.getStringByField("content"); Metadata metadata = (Metadata) tuple.getValueByField("metadata"); // 插入数据到数据库 String sql = "INSERT INTO crawled_pages (url, content, title, crawl_time) VALUES (?, ?, ?, ?)"; try (PreparedStatement pstmt = dbConn.prepareStatement(sql)) { pstmt.setString(1, url); pstmt.setString(2, content); pstmt.setString(3, metadata.getFirstValue("title")); pstmt.setTimestamp(4, new java.sql.Timestamp(System.currentTimeMillis())); pstmt.executeUpdate(); } catch (SQLException e) { e.printStackTrace(); } } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 此Bolt无需输出数据,留空即可 } @Override public void cleanup() { try { if (dbConn != null && !dbConn.isClosed()) { dbConn.close(); } } catch (SQLException e) { e.printStackTrace(); } } }
3. 修改拓扑与配置
在你的 Storm Crawler 拓扑代码中,将自定义 Bolt 连接到解析 Bolt 之后:
TopologyBuilder builder = new TopologyBuilder(); // 配置爬取Spout、Fetch Bolt、Parse Bolt等核心组件 builder.setBolt("parse", new ParseBolt()).shuffleGrouping("fetch"); // 添加自定义数据库存储Bolt builder.setBolt("db-store", new DatabaseStorageBolt()).shuffleGrouping("parse");
同时在项目的crawler-conf.yaml中添加数据库配置:
db.url: "jdbc:mysql://localhost:3306/crawler_db" db.user: "你的数据库用户名" db.password: "你的数据库密码"
4. 编译运行
在项目pom.xml中添加对应数据库的 JDBC 驱动依赖(以 MySQL 为例):
<dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>8.0.33</version> </dependency>
打包成 jar 包后,用 Storm 命令提交拓扑:
storm jar 你的爬虫jar包名.jar 你的拓扑类全路径
三、适配 Typesense/Flex Search
数据存入数据库后,可通过简单脚本定期读取并导入搜索工具,以下是 Typesense 的 Python 示例:
import typesense import mysql.connector # 连接数据库 db_conn = mysql.connector.connect( host="localhost", user="你的数据库用户名", password="你的数据库密码", database="crawler_db" ) cursor = db_conn.cursor(dictionary=True) # 连接Typesense client = typesense.Client({ 'nodes': [{ 'host': '你的Typesense服务器地址', 'port': '8108', 'protocol': 'http' }], 'api_key': '你的Typesense API密钥', 'connection_timeout_seconds': 2 }) # 创建Typesense集合(若不存在) try: client.collections.create({ 'name': 'crawled_pages', 'fields': [ {'name': 'url', 'type': 'string'}, {'name': 'content', 'type': 'string'}, {'name': 'title', 'type': 'string'} ] }) except Exception as e: print("集合已存在:", e) # 读取数据库数据并导入 cursor.execute("SELECT url, content, title FROM crawled_pages") for row in cursor.fetchall(): client.collections['crawled_pages'].documents.create(row) cursor.close() db_conn.close()
四、新手注意事项
- 先做小规模爬取测试,确认数据能正确写入数据库后再扩大爬取范围
- 给数据库的
url字段添加唯一索引,避免重复存储同一页面 - 定期清理无效数据或分表存储,防止数据库数据量过大影响性能
内容的提问来源于stack exchange,提问作者abls1
相关产品推荐
相关产品推荐

