1. 项目概述:从数据洪流中淘金
做量化分析或者基本面研究的朋友,手里没点资金流向数据,心里总是不踏实。股价的涨跌背后,是资金的进进出出在推动。你能看到K线图,能看到成交量,但更细粒度的——主力资金、散户资金、超大单、大单、中单、小单——这些分档数据,才是洞察市场情绪和主力意图的关键。东方财富网的资金流向页面,几乎是国内股民和数据分析师的“数据水源地”之一,它提供了个股和板块非常详尽的历史资金流向明细。
但这个页面是给人看的,不是给程序直接吃的。手动一天天去翻、去抄,效率低到令人发指,而且容易出错。所以,用Python写个爬虫,把这些结构化的资金流向数据自动抓下来,并规整地存入数据库,就成了一个非常实际且高频的需求。这不仅仅是“爬数据”,更是构建个人量化分析数据底座的起点。有了这个自动化流程,你才能腾出手来,专注于更重要的策略研究和模型构建。
接下来,我会带你完整走一遍这个流程。从分析页面结构、设计数据表,到编写稳健的爬虫、处理反爬策略,最后将数据持久化到数据库。我会分享我在这个过程中踩过的坑和总结的技巧,目标是让你拿到一个可以直接运行、易于扩展的解决方案。
2. 核心思路与工具选型
2.1 目标数据源分析与策略制定
东方财富的资金流向数据页面,通常URL模式类似http://data.eastmoney.com/zjlx/股票代码.html。例如,贵州茅台就是http://data.eastmoney.com/zjlx/600519.html。我们的核心目标是抓取页面中的“历史资金流向”表格数据。
策略选择:动态渲染 vs. 静态接口
首先需要判断页面数据是直接加载在HTML中,还是通过JavaScript动态渲染的。这是爬虫设计的第一步,方向错了后面全白费。
- 初步探查:用浏览器打开页面,右键“查看网页源代码”。在源代码里搜索表格里的典型数据,比如“净流入”、“主力资金”。如果搜不到,基本可以断定是动态加载的。
- 确认动态加载:在浏览器开发者工具的“网络”(Network)选项卡中,刷新页面,筛选XHR或Fetch请求。你会很快发现一个关键的接口请求,其URL通常包含
API、GetData等字样,返回的是JSON格式的数据。这才是数据的“真身”。 - 接口分析:找到这个接口后,重点分析它的请求参数。通常包括:
code:股票代码,可能带市场前缀(如SH600519)。type:数据类型,可能区分日、周、月。page和size:分页参数。- 一些时间戳或加密参数,用于反爬。
我们的选择:直接请求这个后端JSON接口,而不是用Selenium或Playwright去模拟浏览器渲染整个页面。理由很充分:效率极高(省去了加载CSS、JS、图片的时间),消耗资源少,速度更快,也更稳定。难点在于如何找到并模拟这个接口的请求。
工具选型背后的逻辑:
- Requests + json:用于发起HTTP请求和解析JSON数据,这是核心。
- Pandas:并非必须,但极其推荐。因为接口返回的JSON数据很容易被转换为DataFrame,而Pandas在数据清洗、转换(如日期格式处理)方面是“神器”,一行代码能顶十行原生Python操作。
- SQLAlchemy:数据库ORM工具。为什么不用简单的
sqlite3模块?因为SQLAlchemy提供了数据库无关的抽象层。今天你用SQLite做测试,明天想迁移到MySQL或PostgreSQL,几乎不需要修改代码。它还能很好地处理连接池、数据类型映射,让代码更健壮、更专业。 - Schedule / APScheduler:如果你需要定时任务(比如每天收盘后自动爬取),这些库能帮你实现可靠的定时调度。
2.2 数据库表结构设计
在写爬虫之前,先想好数据怎么存。一个好的表结构是后续所有分析的基础。资金流向表的核心字段通常包括:
| 字段名 | 数据类型 | 说明 | 设计理由 |
|---|---|---|---|
id | INTEGER PRIMARY KEY | 自增主键 | 每条记录的唯一标识,便于管理和关联。 |
symbol | VARCHAR(10) | 股票代码 | 如600519, 必须建立索引,因为这是最常用的查询条件。 |
trade_date | DATE | 交易日期 | 如2023-10-27, 与symbol组成唯一约束,防止重复插入。 |
close_price | DECIMAL(10,2) | 收盘价 | 记录当日的价格,用于计算和分析。 |
change_pct | DECIMAL(6,2) | 涨跌幅(%) | 股价变动比例。 |
main_net_inflow | DECIMAL(16,2) | 主力净流入(元) | 核心指标,通常指超大单+大单的净额。 |
main_net_inflow_rate | DECIMAL(6,2) | 主力净流入占比(%) | 净流入占成交额的比例,比绝对值更有意义。 |
retail_net_inflow | DECIMAL(16,2) | 散户净流入(元) | 通常指小单的净额,与主力资金形成对照。 |
超大单流入,超大单流出 | DECIMAL(16,2) | 分档资金流 | 更细粒度的数据,有助于深度分析。建议分开存储。 |
turnover_rate | DECIMAL(6,2) | 换手率(%) | 市场活跃度指标。 |
created_at | TIMESTAMP | 记录创建时间 | 默认CURRENT_TIMESTAMP,用于追踪数据获取时间。 |
设计要点与避坑指南:
- 唯一约束:一定要在
(symbol, trade_date)上建立唯一约束(UNIQUE CONSTRAINT)。这是防止数据重复的“铁闸”。爬虫运行时,无论是网络重试还是手动执行,都能确保数据库里同一只股票同一天只有一条记录。 - 索引策略:除了唯一约束自带的索引,如果经常按日期范围查询(如“查询某股票2023年所有数据”),可以考虑在
trade_date上单独建立索引。但索引不是越多越好,会影响写入速度。 - 金额字段精度:资金数据可能很大,使用
DECIMAL(16,2)可以安全存储千亿级别的金额并保留两位小数。DECIMAL类型能保证精确计算,避免浮点数误差。 - 日期类型:务必使用数据库的
DATE类型存储trade_date,而不是VARCHAR。这能让你直接使用数据库的日期函数进行高效的区间查询和聚合。
3. 爬虫核心实现与反爬对抗
3.1 逆向分析接口与参数构造
这是最具技术挑战性的一步。以东方财富为例,我们打开开发者工具,定位到那个返回历史资金流JSON数据的请求。
假设我们找到的接口是:http://datacenter.eastmoney.com/api/data/get?type=RPTA_WEB_RZRQ_GGMX&sty=ALL&source=WEB&...这只是一个示例,实际接口可能会变。
关键步骤:
- 复制为cURL:在Network中找到该请求,右键选择“Copy” -> “Copy as cURL (bash)”。这将给你一个完整的命令行请求。
- 转换与分析:将cURL命令粘贴到在线工具(如 https://curlconverter.com/python/)中,可以将其转换为Python的Requests代码。这能快速得到完整的请求头(headers)和参数(params)。
- 参数解构:仔细研究这些参数。你会发现一些固定参数,如
type、sty、source。以及一些动态参数:code:可能是SH600519或SZ000001这种格式。page:页码。size:每页条数,可能是20、50、100。filter:可能是一个包含查询条件的字符串,如(TRADE_DATE>‘2022-01-01’)。- 最关键的:
t或_后面跟的一串数字,这通常是一个当前时间戳,用于防止缓存。这是最常见的反爬手段之一。
- 请求头:重点关注
User-Agent、Referer、Accept。Referer通常需要设置为资金流向页面的URL,这是服务器验证请求来源的常用方法。
实操心得:
东方财富的接口参数命名可能不那么直观,
filter参数里的表达式语法需要仔细从接口示例中揣摩。一个技巧是,先通过浏览器正常访问,抓取第一页的请求,然后尝试只修改page和size参数来获取更多数据。如果失败,再检查是否有其他动态生成的令牌(token)。
3.2 稳健的爬虫代码编写
下面是一个融合了关键技巧的爬虫核心函数示例:
import requests import pandas as pd import time from datetime import datetime, timedelta import json def fetch_stock_capital_flow(symbol, start_date=None, end_date=None): """ 获取单只股票的历史资金流向数据 Args: symbol (str): 股票代码,如 '600519' start_date (str): 开始日期,格式 '2023-01-01',默认为一年前 end_date (str): 结束日期,格式 '2023-10-27',默认为今天 Returns: pandas.DataFrame: 资金流向数据,列名已清洗 """ # 1. 构造请求头(Headers) - 模拟浏览器 headers = { 'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/119.0.0.0 Safari/537.36', 'Referer': f'http://data.eastmoney.com/zjlx/{symbol}.html', # 关键!Referer必须正确 'Accept': 'application/json, text/plain, */*', 'Accept-Language': 'zh-CN,zh;q=0.9,en;q=0.8', } # 2. 构造请求参数(Params) # 这里需要根据实际找到的接口调整参数名和值 params = { 'type': 'RPTA_WEB_RZRQ_GGMX', # 示例类型,需替换 'sty': 'ALL', 'source': 'WEB', 'code': f'SH{symbol}' if symbol.startswith(('6', '5', '9')) else f'SZ{symbol}', # 判断市场 'page': 1, 'size': 100, # 每页条数,可调大以减少请求次数 'filter': f"(TRADE_DATE>='{start_date}')" if start_date else None, 't': int(time.time() * 1000), # 关键!时间戳防缓存 } # 清理空参数 params = {k: v for k, v in params.items() if v is not None} all_data = [] max_pages = 10 # 安全限制,防止意外死循环 for page in range(1, max_pages + 1): params['page'] = page print(f"正在抓取 {symbol} 第 {page} 页...") try: # 3. 发送请求 # 接口URL需替换为实际地址 response = requests.get( 'http://datacenter.eastmoney.com/api/data/get', params=params, headers=headers, timeout=15 # 设置超时,避免长时间等待 ) response.raise_for_status() # 如果状态码不是200,抛出异常 # 4. 解析JSON json_data = response.json() # 东方财富的接口数据通常嵌套在 `result` -> `data` 路径下 data_list = json_data.get('result', {}).get('data', []) if not data_list: print(f"第 {page} 页无数据,停止爬取。") break all_data.extend(data_list) # 5. 礼貌性延迟,降低请求频率 time.sleep(1.5) # 重要!避免请求过快被封IP except requests.exceptions.RequestException as e: print(f"请求第 {page} 页时发生错误: {e}") break except json.JSONDecodeError as e: print(f"解析第 {page} 页JSON时发生错误: {e}") print(f"响应文本: {response.text[:500]}") # 打印前500字符辅助调试 break # 6. 转换为DataFrame并清洗 if not all_data: return pd.DataFrame() df = pd.DataFrame(all_data) # 重命名列,使其更易懂 column_mapping = { 'SECURITY_CODE': 'symbol', 'TRADE_DATE': 'trade_date', 'CLOSE_PRICE': 'close_price', 'CHANGE_PCT': 'change_pct', 'MAIN_NET_INFLOW': 'main_net_inflow', 'MAIN_NET_INFLOW_RATE': 'main_net_inflow_rate', # ... 其他字段映射 } df.rename(columns=column_mapping, inplace=True) # 转换日期格式 if 'trade_date' in df.columns: df['trade_date'] = pd.to_datetime(df['trade_date']).dt.date # 转换数值类型 numeric_columns = ['close_price', 'change_pct', 'main_net_inflow', 'main_net_inflow_rate'] for col in numeric_columns: if col in df.columns: df[col] = pd.to_numeric(df[col], errors='coerce') # 错误值转为NaN return df注意事项:
- 错误处理:网络请求必须包含
try...except,并处理超时、状态码异常、JSON解析失败等情况。爬虫要能“优雅地失败”,而不是整个崩溃。 - 延迟(Sleep):
time.sleep(random.uniform(1, 3))比固定延迟更好,模拟人类操作的不确定性。这是对目标网站最基本的尊重,也是保护自己IP不被封禁的有效手段。 - User-Agent轮换:如果需要大规模爬取,准备一个User-Agent列表进行轮换,可以进一步降低被识别风险。
4. 数据持久化:写入数据库
数据抓取到Pandas DataFrame后,下一步就是高效、安全地存入数据库。这里我们使用SQLAlchemy配合SQLite(本地开发测试首选)或MySQL/PostgreSQL。
4.1 使用SQLAlchemy建立连接与映射
from sqlalchemy import create_engine, Column, Integer, String, Date, DECIMAL, TIMESTAMP, UniqueConstraint from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from datetime import datetime # 1. 定义基类 Base = declarative_base() # 2. 定义数据表模型(ORM类) class StockCapitalFlow(Base): __tablename__ = 'stock_capital_flow' id = Column(Integer, primary_key=True, autoincrement=True) symbol = Column(String(10), nullable=False, index=True) # 建立索引 trade_date = Column(Date, nullable=False) close_price = Column(DECIMAL(10, 2)) change_pct = Column(DECIMAL(6, 2)) main_net_inflow = Column(DECIMAL(16, 2)) main_net_inflow_rate = Column(DECIMAL(6, 2)) retail_net_inflow = Column(DECIMAL(16, 2)) # ... 其他字段 created_at = Column(TIMESTAMP, default=datetime.now) # 3. 定义唯一约束,防止重复数据 __table_args__ = (UniqueConstraint('symbol', 'trade_date', name='uix_symbol_trade_date'),) # 4. 创建数据库连接引擎 # 使用SQLite(文件数据库,适合本地测试) # engine = create_engine('sqlite:///./stock_data.db', echo=False) # echo=True 可查看SQL日志 # 使用MySQL(生产环境推荐) # 格式:mysql+pymysql://用户名:密码@主机:端口/数据库名 engine = create_engine('mysql+pymysql://user:password@localhost:3306/stock_db?charset=utf8mb4') # 5. 创建表(如果不存在) Base.metadata.create_all(engine) # 6. 创建会话工厂 SessionLocal = sessionmaker(bind=engine)4.2 高效批量写入与去重策略
直接将DataFrame一行行插入效率极低。Pandas的to_sql方法虽然方便,但默认的插入方式(if_exists='append')无法处理唯一约束冲突,会导致整个插入失败。我们需要更精细的控制。
推荐方案:使用SQLAlchemy Core进行批量“upsert”
“Upsert”是Update和Insert的合成词,即存在则更新,不存在则插入。MySQL 5.7+和PostgreSQL都支持ON DUPLICATE KEY UPDATE语法,SQLite也有类似的INSERT OR REPLACE/IGNORE。
def upsert_capital_flow_to_db(df, engine): """ 将DataFrame数据高效写入数据库,并处理重复数据。 Args: df (pd.DataFrame): 清洗后的资金流向数据 engine: SQLAlchemy引擎 """ if df.empty: print("DataFrame为空,跳过写入。") return # 确保列名与数据库模型一致 # 假设df的列名已经是我们定义好的英文名 # 方法:使用Pandas + SQLAlchemy Core执行原生SQL进行批量upsert # 这里以MySQL为例 table_name = 'stock_capital_flow' temp_table_name = f'temp_{table_name}' # 创建一个临时表,用于暂存本次抓取的数据 with engine.begin() as conn: # 使用事务 # 1. 将DataFrame写入临时表 df.to_sql(temp_table_name, conn, if_exists='replace', index=False) # 2. 执行Upsert操作 # 核心SQL语句:将临时表的数据合并到主表,冲突时更新某些字段 upsert_sql = f""" INSERT INTO {table_name} (symbol, trade_date, close_price, change_pct, main_net_inflow, main_net_inflow_rate, ...) SELECT symbol, trade_date, close_price, change_pct, main_net_inflow, main_net_inflow_rate, ... FROM {temp_table_name} ON DUPLICATE KEY UPDATE close_price = VALUES(close_price), change_pct = VALUES(change_pct), main_net_inflow = VALUES(main_net_inflow), main_net_inflow_rate = VALUES(main_net_inflow_rate), -- ... 更新其他需要覆盖的字段,created_at 通常不更新 updated_at = CURRENT_TIMESTAMP; -- 可以加一个更新时间的字段 """ conn.execute(text(upsert_sql)) # 3. 删除临时表 conn.execute(text(f"DROP TABLE IF EXISTS {temp_table_name}")) print(f"成功写入/更新 {len(df)} 条记录到数据库。")实操心得:
批量Upsert是生产级爬虫数据入库的关键。直接循环插入每秒可能只能处理几十条,而批量Upsert每秒可以处理上万条。临时表的思路避免了在Python内存中拼接庞大的SQL语句,性能和安全性都更好。记住,对于像收盘价、资金流这种客观历史数据,一旦写入,即使重复抓取,也应该用新数据覆盖旧数据(如果接口数据有修正的话)。
created_at记录首次抓取时间,updated_at记录最后一次更新时间,这个设计很实用。
5. 项目整合与调度实战
5.1 构建可配置的主程序
将上述模块组装起来,并增加一些工程化配置。
# config.py import os from dataclasses import dataclass @dataclass class Config: # 数据库连接字符串 DATABASE_URL = os.getenv('DATABASE_URL', 'sqlite:///./stock_data.db') # 请求基础延迟(秒) BASE_DELAY = 1.5 # 每页数据大小 PAGE_SIZE = 100 # 要爬取的股票代码列表 STOCK_LIST = ['600519', '000858', '300750', '002594'] # 示例 # 数据抓取时间范围 START_DATE = '2023-01-01' END_DATE = '2023-10-27' # main.py import logging from config import Config from capital_flow_spider import fetch_stock_capital_flow from database import upsert_capital_flow_to_db, engine import time logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') logger = logging.getLogger(__name__) def main(): config = Config() failed_symbols = [] for symbol in config.STOCK_LIST: logger.info(f"开始处理股票: {symbol}") try: # 抓取数据 df = fetch_stock_capital_flow( symbol=symbol, start_date=config.START_DATE, end_date=config.END_DATE ) if df.empty: logger.warning(f"股票 {symbol} 未获取到数据。") continue logger.info(f"股票 {symbol} 抓取到 {len(df)} 条记录。") # 写入数据库 upsert_capital_flow_to_db(df, engine) # 随机延迟,避免对服务器造成压力 time.sleep(config.BASE_DELAY) except Exception as e: logger.error(f"处理股票 {symbol} 时发生错误: {e}", exc_info=True) failed_symbols.append(symbol) # 可以在每个股票处理完后加一个稍长的延迟 time.sleep(2) if failed_symbols: logger.error(f"以下股票处理失败: {failed_symbols}") else: logger.info("所有股票处理完成!") if __name__ == '__main__': main()5.2 定时任务与日志管理
对于需要每日更新的场景,定时任务必不可少。
# scheduler.py import schedule import time from main import main import logging def job(): logging.info("定时爬虫任务开始执行...") try: main() logging.info("定时爬虫任务执行完毕。") except Exception as e: logging.error(f"定时任务执行失败: {e}", exc_info=True) if __name__ == '__main__': # 设置日志 logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', handlers=[ logging.FileHandler('capital_flow_spider.log'), # 输出到文件 logging.StreamHandler() # 输出到控制台 ] ) # 每天下午16:30(收盘后)执行一次 schedule.every().day.at("16:30").do(job) logging.info("爬虫定时调度器已启动,等待执行...") while True: schedule.run_pending() time.sleep(60) # 每分钟检查一次日志管理:将日志同时输出到文件和控制台,便于后期排查问题。日志文件要设置轮转(Rotating),避免单个文件过大。
6. 常见问题与排查技巧实录
在实际运行中,你肯定会遇到各种问题。这里记录了几个最典型的“坑”和解决方法。
6.1 接口变更或数据为空
问题:昨天还能跑通的爬虫,今天突然抓不到数据了,或者返回空列表。排查:
- 手动访问页面:首先用浏览器打开目标页面,确认数据是否正常显示。如果页面都变了,那接口大概率也变了。
- 检查网络请求:再次打开开发者工具,查看之前那个关键的XHR请求是否还在,参数是否有变化。重点看
type、filter格式、是否有新增的必选参数(如token、v版本号)。 - 模拟请求:用Postman或直接在Python脚本里打印出完整的请求URL和响应内容(
response.text),看看服务器返回了什么错误信息(如{“code“: -1, “message“: “参数错误“})。
解决:根据新的接口文档或抓包结果,调整爬虫的请求参数和URL。这是爬虫维护的常态。
6.2 IP被封禁或访问频率限制
问题:请求返回403 Forbidden、429 Too Many Requests,或者返回一些提示“操作频繁”的HTML页面。表现:最初几次请求成功,连续快速请求后开始失败。解决:
- 首要措施:大幅降低请求频率。这是最有效的方法。在循环中增加
time.sleep(random.uniform(3, 7)),让请求间隔更随机、更长。 - 使用代理IP池:如果需要爬取大量股票或历史数据,这是必选项。可以购买付费代理服务,或者自建代理池。在Requests中设置代理很简单:
proxies = {“http“: “http://your-proxy:port“, “https“: “https://your-proxy:port“},然后在请求中传入。 - 完善请求头:确保
User-Agent、Referer、Accept-Language、Cookie(如果需要)等头部信息与真实浏览器一致。有些网站会校验这些信息。
6.3 数据字段映射错误或类型转换异常
问题:数据能抓到,但存入数据库时出错,比如字符串无法转为数字,日期格式解析失败。排查:
- 打印原始数据:在清洗数据前,打印几行原始的JSON数据,查看字段名和值到底是什么样子。东方财富的接口字段名可能是中文,也可能是缩写英文。
- 逐步清洗:不要试图一步到位。先确保能正确解析JSON为列表,再转换为DataFrame,然后一步一步进行重命名、类型转换。在每个步骤后打印
df.dtypes和df.head()来检查。 - 处理缺失值:接口返回的某些字段可能为空字符串
““或null。使用pd.to_numeric(..., errors=‘coerce‘)可以将无效解析转为NaN,避免程序崩溃。
6.4 数据库连接与并发写入问题
问题:在多进程或异步爬虫中,同时写入数据库导致锁表、连接超时或重复数据。解决:
- 连接池:SQLAlchemy的
create_engine默认启用了连接池。确保你的数据库连接字符串是全局唯一的,并在所有爬虫线程/进程中共享这个引擎。 - 会话管理:每个线程或进程使用独立的
Session,并在任务结束后正确关闭 (session.close())。不要跨线程共享Session。 - 更细粒度的锁:如果使用SQLite,写入密集型操作建议用
WAL(Write-Ahead Logging) 模式,可以提高并发性。对于MySQL/PostgreSQL,确保使用InnoDB等支持行级锁的引擎。 - 终极方案:消息队列:对于超大规模爬虫,可以将爬取任务和数据写入解耦。爬虫进程只负责抓取,将数据扔到Redis或RabbitMQ这样的消息队列中,再由独立的消费者进程从队列中取出数据,批量写入数据库。这能极大提升系统的稳定性和扩展性。
6.5 数据更新逻辑的权衡
问题:每天跑爬虫,是全部重新抓取覆盖,还是只抓取新增日期?方案:
- 全量覆盖:逻辑简单,代码好写。每天定时任务跑一遍所有历史数据。缺点是浪费带宽和计算资源,且对目标网站不友好。
- 增量更新:更优雅。先从数据库查询某只股票最新的
trade_date,然后只请求这个日期之后的数据。这需要爬虫接口支持按日期过滤。优点是高效、友好。推荐此方案。
# 增量更新示例片段 def get_latest_date_from_db(symbol, session): result = session.query(StockCapitalFlow.trade_date).filter_by(symbol=symbol).order_by(StockCapitalFlow.trade_date.desc()).first() return result[0] if result else None # 在抓取前调用 latest_date = get_latest_date_from_db(symbol, session) start_date = (latest_date + timedelta(days=1)).strftime(‘%Y-%m-%d‘) if latest_date else config.START_DATE # 然后用这个 start_date 去构造请求参数爬虫项目从来不是“一劳永逸”的,它是一个需要持续观察、维护和调整的系统。核心在于构建一个健壮的框架,当数据源发生变化时,你能快速定位问题并修复。把数据抓下来只是第一步,确保数据管道长期稳定、准确、高效地运行,才是真正的价值所在。