You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 的 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 12:15:25