ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Python爬虫实战:抓取东方财富股票资金流向数据并入库MySQL

2026/8/5 5:12:10 拓冰建站 浏览量
Python爬虫实战:抓取东方财富股票资金流向数据并入库MySQL

1. 项目缘起与核心价值

最近在整理自己的量化分析策略时,发现一个挺实际的需求:想回溯分析某只股票在特定时间段内的资金流向变化,看看主力资金的进出节奏和股价波动有没有什么关联。市面上虽然有很多数据接口,但要么收费不菲,要么数据维度不全,特别是像东方财富网提供的这种按日统计的“历史资金流向”明细表——包括主力净流入、超大单、大单、中单、小单的净额数据——在很多免费API里是缺失的。手动去网站复制粘贴?对于多只股票、长时间跨度的分析来说,这工作量简直是个噩梦。于是,一个念头自然就冒出来了:写个Python爬虫,自动抓取东方财富网上的股票历史资金流向数据,然后规规矩矩地存进自己的数据库里,方便后续做任何分析和回测。

这个项目听起来就是“爬虫+数据存储”,但里面有几个关键点值得深究。首先,东方财富作为主流财经门户,其反爬机制一直在升级,直接requests.get大概率会吃闭门羹。其次,资金流向数据通常是通过异步加载(Ajax)动态生成的,网址(URL)看起来可能很规整,但里面往往藏着校验参数。最后,数据入库不是简单存个CSV文件就完事了,需要考虑表结构设计、增量更新、避免重复存储以及后续查询效率等问题。把这些环节打通,你得到的不仅是一份数据,更是一套可复用、可扩展的数据获取与管理系统的基础框架。无论你是量化交易新手想积累研究素材,还是数据分析师需要稳定的数据源,这个项目都能提供一个扎实的起点。

2. 逆向解析:定位东方财富资金流向的真实数据接口

动手写代码之前,最关键的一步是找到数据从哪里来。很多新手会直接去爬取东方财富股票详情页的HTML,比如quote.eastmoney.com/sh600036.html这样的页面。但如果你打开开发者工具(F12),切换到“网络”(Network)选项卡,然后点击页面上的“资金流向”或“历史资金流向”标签,你会发现浏览器发出了新的请求。

2.1 寻找核心API请求

以浦发银行(600036)为例,在历史资金流向页面,通过筛选XHR/Fetch请求,你很可能会发现一个类似以下的请求:https://datacenter.eastmoney.com/securities/api/data/get?type=RPTA_WEB_SUPERVISOR_MONEY&sty=ALL&source=WEB&client=WEB&filter=(SECURITY_CODE%3D%22600036%22)&p=1&ps=5000&sr=-1&st=TRADE_DATE&var=xxxxxx

这个URL包含了几个重要信息:

  • type=RPTA_WEB_SUPERVISOR_MONEY: 这很可能就是获取监管资金流向数据的接口类型。
  • filter=(SECURITY_CODE%3D%22600036%22): 这是URL编码后的过滤条件,指定了股票代码SECURITY_CODE="600036"
  • p=1ps=5000: 代表页码和每页大小,这里请求第一页,每页5000条,基本可以一次性拿到所有历史数据。
  • st=TRADE_DATEsr=-1: 可能代表排序字段和顺序(按交易日期倒序)。

2.2 参数分析与简化

实际测试中,一些参数可能是非必需的或者可以固定。我们可以尝试简化请求。最关键的部分通常是typefilter和分页参数。var参数看起来像是一个随机数或时间戳,用于防止缓存,但在requests库中,我们通常可以通过添加时间戳参数或禁用缓存来达到类似效果。

经过测试和验证,一个稳定可用的请求URL可能简化为:https://datacenter.eastmoney.com/securities/api/data/get?type=RPTA_WEB_SUPERVISOR_MONEY&sty=ALL&source=WEB&client=WEB&filter=(SECURITY_CODE%3D“股票代码”)&p=1&ps=5000&st=TRADE_DATE&sr=-1

这里的核心是构造filter条件。对于不同的股票,你只需要替换SECURITY_CODE的值即可。对于沪深主板,代码就是6位数字,如“600036”;对于创业板,可能是“300750”;对于科创板,是“688981”。注意,接口可能要求代码带市场前缀,如“SH600036”“SZ000001”,这需要通过实际测试和观察返回数据来确定。

2.3 解析返回的JSON数据结构

这个接口返回的数据通常是JSON格式。成功请求后,你会得到一个结构复杂的JSON对象。你需要层层剥开,找到真正的数据列表。路径可能类似于data -> result -> data

数据列表中的每一条,通常对应一个交易日的资金流向汇总信息,包含的字段可能有:

  • TRADE_DATE: 交易日期
  • SECURITY_CODE: 股票代码
  • SECURITY_NAME: 股票名称
  • MAIN_NET_INFLOW: 主力净流入(元)
  • MAIN_NET_INFLOW_RATIO: 主力净流入占比(%)
  • HUGE_ORDER_NET_INFLOW: 超大单净流入(元)
  • BIG_ORDER_NET_INFLOW: 大单净流入(元)
  • MID_ORDER_NET_INFLOW: 中单净流入(元)
  • SMALL_ORDER_NET_INFLOW: 小单净流入(元)
  • CLOSE_PRICE: 收盘价
  • CHANGE_RATE: 涨跌幅(%)

注意:字段名(Key)可能因接口版本而变化。务必在首次成功获取数据后,仔细打印并查看完整的JSON响应结构,确认上述字段的确切名称。这是避免后续数据处理出错的基础。

3. 构建稳健的Python爬虫:从请求到数据清洗

找到了接口,接下来就是用Python把它自动化。这里我们选择requests库进行网络请求,用pandas进行便捷的数据处理。

3.1 环境准备与依赖安装

首先确保你的Python环境已经安装了必要的库。可以通过pip安装:

pip install requests pandas sqlalchemy pymysql
  • requests: 用于发送HTTP请求。
  • pandas: 数据处理和分析的核心库,能轻松将JSON数据转为DataFrame。
  • sqlalchemy: Python的SQL工具包和对象关系映射(ORM)工具,这里我们主要用它的create_engine来建立数据库连接,写法更通用。
  • pymysql: MySQL数据库的Python驱动。如果你用PostgreSQL或SQLite,则对应安装psycopg2sqlite3(后者通常内置)。

3.2 构造请求头与应对反爬

东方财富的接口对请求头有一定检查。直接使用requests.get()而不设置头信息,可能会被拒绝或返回错误数据。一个相对安全的做法是模拟常见浏览器的请求头。

import requests import pandas as pd from datetime import datetime def fetch_money_flow(stock_code): """ 获取单只股票的历史资金流向数据 :param stock_code: 股票代码,如 '600036' (需确认是否带市场前缀) :return: 包含资金流向数据的pandas DataFrame, 失败则返回空DataFrame """ # 构造请求URL url = “https://datacenter.eastmoney.com/securities/api/data/get” params = { ‘type’: ‘RPTA_WEB_SUPERVISOR_MONEY’, ‘sty’: ‘ALL’, ‘source’: ‘WEB’, ‘client’: ‘WEB’, ‘filter’: f’(SECURITY_CODE=“{stock_code}”)’, # 注意这里的引号是英文的 ‘p’: 1, ‘ps’: 5000, ‘st’: ‘TRADE_DATE’, ‘sr’: ‘-1’, } # 模拟浏览器请求头 headers = { ‘User-Agent’: ‘Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36’, ‘Accept’: ‘application/json, text/plain, */*’, ‘Accept-Language’: ‘zh-CN,zh;q=0.9,en;q=0.8’, ‘Origin’: ‘https://data.eastmoney.com’, ‘Referer’: f’https://data.eastmoney.com/zjlx/{stock_code}.html’, # 这个Referer很重要 ‘Sec-Fetch-Dest’: ‘empty’, ‘Sec-Fetch-Mode’: ‘cors’, ‘Sec-Fetch-Site’: ‘same-site’, } try: response = requests.get(url, params=params, headers=headers, timeout=10) response.raise_for_status() # 检查HTTP请求是否成功 json_data = response.json() # 解析JSON数据,路径需要根据实际返回结构调整 # 常见路径:json_data[‘data’][‘result’][‘data’] data_list = json_data.get(‘data’, {}).get(‘result’, {}).get(‘data’, []) if not data_list: print(f“未获取到股票 {stock_code} 的资金流向数据。”) return pd.DataFrame() # 将数据列表转换为DataFrame df = pd.DataFrame(data_list) return df except requests.exceptions.RequestException as e: print(f“请求股票 {stock_code} 数据时发生错误: {e}”) return pd.DataFrame() except ValueError as e: print(f“解析股票 {stock_code} 的JSON数据时发生错误: {e}”) return pd.DataFrame()

3.3 数据清洗与格式化

从接口拿到的原始DataFrame可能包含不需要的字段,或者字段格式不理想(比如日期是时间戳或特定字符串)。我们需要进行清洗。

def clean_money_flow_data(df, stock_code): """ 清洗和格式化资金流向数据 :param df: 原始的DataFrame :param stock_code: 股票代码,用于填充缺失 :return: 清洗后的DataFrame """ if df.empty: return df # 1. 选择需要的列,并重命名为更易读的英文名 # 注意:列名必须与JSON中的key完全一致,这里仅为示例 column_mapping = { ‘TRADE_DATE’: ‘trade_date’, ‘SECURITY_CODE’: ‘symbol’, ‘SECURITY_NAME’: ‘name’, ‘MAIN_NET_INFLOW’: ‘main_net_inflow’, ‘MAIN_NET_INFLOW_RATIO’: ‘main_net_inflow_ratio’, ‘HUGE_ORDER_NET_INFLOW’: ‘huge_order_net’, ‘BIG_ORDER_NET_INFLOW’: ‘big_order_net’, ‘MID_ORDER_NET_INFLOW’: ‘mid_order_net’, ‘SMALL_ORDER_NET_INFLOW’: ‘small_order_net’, ‘CLOSE_PRICE’: ‘close’, ‘CHANGE_RATE’: ‘change_pct’, } # 只保留我们映射表中存在的列 existing_columns = [col for col in column_mapping.keys() if col in df.columns] df = df[existing_columns].copy() df.rename(columns=column_mapping, inplace=True) # 2. 处理日期字段:假设原始是‘2023-04-28 00:00:00’这样的字符串 if ‘trade_date’ in df.columns: df[‘trade_date’] = pd.to_datetime(df[‘trade_date’]).dt.date # 只保留日期部分 # 3. 确保股票代码列存在且一致 if ‘symbol’ not in df.columns or df[‘symbol’].isnull().all(): df[‘symbol’] = stock_code # 4. 处理数值字段:去除逗号,转换为浮点数(如果原始数据是字符串的话) numeric_columns = [‘main_net_inflow’, ‘huge_order_net’, ‘big_order_net’, ‘mid_order_net’, ‘small_order_net’, ‘close’, ‘main_net_inflow_ratio’, ‘change_pct’] for col in numeric_columns: if col in df.columns: # 先转换为字符串,再替换非数字字符(如逗号、百分号) df[col] = df[col].astype(str).str.replace(‘,’, ‘’).str.replace(‘%’, ‘’) # 转换为浮点数,错误强制转为NaN df[col] = pd.to_numeric(df[col], errors=‘coerce’) # 5. 按交易日期排序(通常接口已返回排序数据,这里确保一下) if ‘trade_date’ in df.columns: df.sort_values(by=‘trade_date’, ascending=False, inplace=True) # 按日期倒序,最新在前 df.reset_index(drop=True, inplace=True) # 6. 添加数据获取时间戳 df[‘created_at’] = datetime.now() return df

实操心得:在数据清洗时,一定要先打印出几行原始数据看看格式。东方财富的数值字段有时会是以“万”或“亿”为单位的字符串,或者带有百分号。pd.to_numericerrors=‘coerce’参数非常有用,它会把转换失败的值变成NaN,而不是让整个程序崩溃,方便后续排查。

4. 数据库设计:为股票资金流向数据安家

数据抓取和清洗完成后,下一步就是持久化存储。选择MySQL作为数据库,主要是因为它普及率高、生态成熟,与Python(通过pymysql+sqlalchemy)和很多数据分析工具(如pandasMetabase)集成起来非常方便。

4.1 表结构设计思路

设计表结构时,要考虑以下几点:

  1. 唯一性约束:同一个股票在同一天的数据不应该重复存储。最自然的唯一键是(symbol, trade_date)
  2. 字段类型
    • trade_date: 使用DATE类型,比DATETIME更节省空间,也符合业务语义。
    • 资金净额字段:如main_net_inflow,单位是“元”,可能数值很大,使用DECIMAL(20, 2)BIGINT(以分为单位存储)都可以。DECIMAL(20,2)能直接存储元为单位、保留两位小数的数值,更直观。
    • 比例字段:如main_net_inflow_ratio,使用DECIMAL(8, 4)足够,表示-9999.9999%9999.9999%的范围。
    • 股票代码和名称:使用VARCHAR
  3. 索引:为了加快按股票代码和日期范围的查询速度,必须在(symbol, trade_date)上建立复合索引,并且由于它是唯一键,本身就带有索引。如果经常需要按单个字段查询,可以考虑单独为trade_date建立索引。
  4. 元信息:添加created_at字段记录数据插入时间,用于审计和追踪。

4.2 创建数据表的SQL语句

CREATE TABLE stock_money_flow ( id INT AUTO_INCREMENT PRIMARY KEY COMMENT ‘自增主键’, symbol VARCHAR(10) NOT NULL COMMENT ‘股票代码,如 600036.SH’, trade_date DATE NOT NULL COMMENT ‘交易日期’, name VARCHAR(50) COMMENT ‘股票名称’, main_net_inflow DECIMAL(20, 2) COMMENT ‘主力净流入(元)’, main_net_inflow_ratio DECIMAL(8, 4) COMMENT ‘主力净流入占比(%)’, huge_order_net DECIMAL(20, 2) COMMENT ‘超大单净流入(元)’, big_order_net DECIMAL(20, 2) COMMENT ‘大单净流入(元)’, mid_order_net DECIMAL(20, 2) COMMENT ‘中单净流入(元)’, small_order_net DECIMAL(20, 2) COMMENT ‘小单净流入(元)’, close DECIMAL(10, 2) COMMENT ‘收盘价’, change_pct DECIMAL(8, 4) COMMENT ‘涨跌幅(%)’, created_at DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT ‘数据创建时间’, UNIQUE KEY uk_symbol_date (symbol, trade_date), -- 唯一约束,防止重复 KEY idx_trade_date (trade_date) -- 按日期查询的索引 ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT=‘股票历史资金流向表’;

4.3 使用SQLAlchemy连接与操作数据库

在Python中,我们使用SQLAlchemy的create_engine来创建数据库连接,它提供了一个统一的接口,即使以后换数据库(比如PostgreSQL或SQLite),代码改动也很小。

from sqlalchemy import create_engine, text from urllib.parse import quote_plus def create_db_connection(db_config): """ 创建数据库连接引擎 :param db_config: 字典,包含host, port, user, password, database等信息 :return: sqlalchemy.engine.Engine 实例 """ # 对密码进行URL编码,防止特殊字符导致连接失败 encoded_password = quote_plus(db_config[‘password’]) # 构建连接字符串 # 格式: dialect+driver://username:password@host:port/database connection_string = f“mysql+pymysql://{db_config[‘user’]}:{encoded_password}@{db_config[‘host’]}:{db_config[‘port’]}/{db_config[‘database’]}” # 创建引擎,echo=True可以在控制台看到执行的SQL,调试时有用 engine = create_engine(connection_string, echo=False) return engine

5. 数据入库策略:增量更新与避免重复

最直接的入库方式是把抓取到的所有数据一次性插入(df.to_sql(..., if_exists=‘append’))。但这会带来两个问题:1) 如果程序多次运行,会导致数据重复;2) 每次都是全量插入,效率低下。我们需要实现增量更新。

5.1 “存在即更新,不存在则插入”策略

在MySQL中,我们可以使用INSERT ... ON DUPLICATE KEY UPDATE语句。这正好利用了我们在表设计中设置的唯一键(symbol, trade_date)。当插入的数据与已有唯一键冲突时,就执行更新操作。

Pandas的to_sql方法原生不支持这个语法,但我们可以通过一些方法实现。

方法一:使用SQLAlchemy Core逐条处理(清晰但稍慢)

def save_to_db_incremental(df, engine, table_name=‘stock_money_flow’): """ 使用增量方式(INSERT ... ON DUPLICATE KEY UPDATE)保存数据到数据库 :param df: 清洗后的DataFrame :param engine: SQLAlchemy引擎 :param table_name: 表名 """ if df.empty: print(“没有数据需要保存。”) return # 获取数据库连接 with engine.connect() as connection: # 开始一个事务 with connection.begin(): for index, row in df.iterrows(): # 构建INSERT ... ON DUPLICATE KEY UPDATE语句 # 这里假设所有字段都需要在冲突时更新,除了唯一键和自增主键 placeholders = ‘, ‘.join([f‘:{col}’ for col in df.columns]) columns = ‘, ‘.join(df.columns) update_clause = ‘, ‘.join([f‘{col}=VALUES({col})’ for col in df.columns if col not in [‘id’, ‘symbol’, ‘trade_date’]]) sql = text(f“”” INSERT INTO {table_name} ({columns}) VALUES ({placeholders}) ON DUPLICATE KEY UPDATE {update_clause} “””) # 将行转换为字典并执行 params = row.to_dict() connection.execute(sql, params) print(f“成功增量更新 {len(df)} 条记录到表 {table_name}。”)

方法二:使用pandas配合临时表(高效,适合大批量)这种方法更高效,尤其当数据量较大时。思路是:1) 将DataFrame写入一个临时表;2) 用一条SQL语句将临时表的数据合并到主表。

def save_to_db_incremental_bulk(df, engine, table_name=‘stock_money_flow’, temp_table_name=‘temp_money_flow’): """ 使用临时表进行批量增量更新(推荐用于大批量数据) :param df: 清洗后的DataFrame :param engine: SQLAlchemy引擎 :param table_name: 目标表名 :param temp_table_name: 临时表名 """ if df.empty: return with engine.connect() as connection: with connection.begin(): # 1. 将DataFrame写入临时表(每次覆盖) df.to_sql(temp_table_name, con=engine, if_exists=‘replace’, index=False) # 2. 执行合并操作 # 构建列名列表 columns = df.columns.tolist() columns_str = ‘, ‘.join(columns) update_assignments = ‘, ‘.join([f‘{col}=t.{col}’ for col in columns if col not in [‘id’, ‘symbol’, ‘trade_date’]]) merge_sql = text(f“”” INSERT INTO {table_name} ({columns_str}) SELECT {columns_str} FROM {temp_table_name} t ON DUPLICATE KEY UPDATE {update_assignments} “””) connection.execute(merge_sql) # 3. 删除临时表(可选) # connection.execute(text(f“DROP TABLE IF EXISTS {temp_table_name}”)) print(f“通过临时表批量更新了 {len(df)} 条记录。”)

踩坑提醒:使用“方法二”时,临时表的结构必须与目标表完全一致(列名、顺序、类型)。df.to_sql创建临时表时可能会自动推断类型,可能与目标表有细微差异(比如字符串长度),可能导致合并失败。一个更稳妥的做法是先用CREATE TABLE ... LIKE ...语句创建一个结构相同的临时表,然后再插入数据。

5.2 主程序流程整合

现在,我们把所有模块串联起来,形成一个完整的脚本。

import time from sqlalchemy.exc import SQLAlchemyError def main(): # 1. 数据库配置 db_config = { ‘host’: ‘localhost’, ‘port’: 3306, ‘user’: ‘your_username’, ‘password’: ‘your_password’, ‘database’: ‘stock_data’ } # 2. 股票代码列表 stock_symbols = [‘600036’, ‘000001’, ‘300750’] # 示例代码 # 3. 创建数据库连接 try: engine = create_db_connection(db_config) except Exception as e: print(f“数据库连接失败: {e}”) return # 4. 遍历股票代码,抓取并保存数据 for symbol in stock_symbols: print(f“正在处理股票: {symbol}”) # 4.1 抓取数据 raw_df = fetch_money_flow(symbol) if raw_df.empty: print(f“ -> 股票 {symbol} 数据抓取失败或为空,跳过。”) continue # 4.2 清洗数据 cleaned_df = clean_money_flow_data(raw_df, symbol) if cleaned_df.empty: print(f“ -> 股票 {symbol} 数据清洗后为空,跳过。”) continue print(f“ -> 成功获取 {len(cleaned_df)} 条记录。”) # 4.3 保存到数据库 try: # 选择一种增量更新方法 # save_to_db_incremental(cleaned_df, engine) # 方法一 save_to_db_incremental_bulk(cleaned_df, engine) # 方法二(推荐) except SQLAlchemyError as e: print(f“ -> 股票 {symbol} 数据入库失败: {e}”) # 4.4 礼貌延时,避免请求过快被封IP time.sleep(2) print(“所有股票数据处理完毕。”) engine.dispose() # 关闭引擎 if __name__ == ‘__main__’: main()

6. 项目优化与进阶思考

一个能跑通的脚本只是开始。要让这个数据管道真正可靠、高效、可维护,还需要考虑更多。

6.1 错误处理与重试机制

网络请求和数据库操作都可能失败。必须添加健壮的错误处理和重试逻辑。对于网络请求,可以使用tenacityretrying库实现指数退避重试。

from tenacity import retry, stop_after_attempt, wait_exponential @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) def fetch_money_flow_with_retry(stock_code): """带重试的抓取函数""" return fetch_money_flow(stock_code) # 调用我们之前定义的函数

6.2 日志记录

print输出信息不利于后期排查问题。应该使用Python的logging模块,将不同级别的信息(INFO, WARNING, ERROR)输出到文件和控制台。

6.3 配置化管理

数据库连接信息、请求头、目标股票列表等不应该硬编码在脚本里。可以使用配置文件(如config.iniconfig.yaml)或环境变量来管理。

6.4 定时任务与自动化

要让数据每天自动更新,可以结合操作系统的定时任务(如Linux的cron或Windows的任务计划程序)来定期执行这个Python脚本。更复杂的调度可以使用AirflowPrefect这样的工作流管理工具。

6.5 数据验证与监控

数据入库后,可以添加简单的验证步骤,比如检查最新日期的数据是否已成功入库,或者统计每日入库的数据量是否在合理范围内。这可以通过在脚本最后执行一些查询语句来实现。

6.6 扩展性考虑

  • 更多数据源:东方财富的接口可能不稳定或有限制。可以同时集成其他免费数据源(如AKShare库,它封装了多个财经数据接口),作为备份或补充。
  • 更多数据维度:除了资金流向,还可以并行抓取日K线、财务指标、龙虎榜等数据,丰富你的数据库。
  • 数据质量:定期检查数据中的异常值(如涨跌幅超过限制、资金流数据为0的交易日等),并建立清洗规则。

这个项目从简单的需求出发,串联起了网络爬虫、数据清洗、数据库操作、错误处理等多个核心技能点。把它跑通并不断优化,你收获的不仅仅是一堆股票数据,更是一套处理数据管道问题的实战经验。在实际操作中,最大的挑战往往不是代码本身,而是目标网站的反爬策略变化,这就需要你保持对网络请求的敏锐观察,随时准备调整你的爬虫策略。