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?typeRPTA_WEB_SUPERVISOR_MONEYstyALLsourceWEBclientWEBfilter(SECURITY_CODE%3D%22600036%22)p1ps5000sr-1stTRADE_DATEvarxxxxxx这个URL包含了几个重要信息typeRPTA_WEB_SUPERVISOR_MONEY: 这很可能就是获取监管资金流向数据的接口类型。filter(SECURITY_CODE%3D%22600036%22): 这是URL编码后的过滤条件指定了股票代码SECURITY_CODE600036。p1和ps5000: 代表页码和每页大小这里请求第一页每页5000条基本可以一次性拿到所有历史数据。stTRADE_DATE和sr-1: 可能代表排序字段和顺序按交易日期倒序。2.2 参数分析与简化实际测试中一些参数可能是非必需的或者可以固定。我们可以尝试简化请求。最关键的部分通常是type、filter和分页参数。var参数看起来像是一个随机数或时间戳用于防止缓存但在requests库中我们通常可以通过添加时间戳参数或禁用缓存来达到类似效果。经过测试和验证一个稳定可用的请求URL可能简化为https://datacenter.eastmoney.com/securities/api/data/get?typeRPTA_WEB_SUPERVISOR_MONEYstyALLsourceWEBclientWEBfilter(SECURITY_CODE%3D“股票代码”)p1ps5000stTRADE_DATEsr-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 pymysqlrequests: 用于发送HTTP请求。pandas: 数据处理和分析的核心库能轻松将JSON数据转为DataFrame。sqlalchemy: Python的SQL工具包和对象关系映射ORM工具这里我们主要用它的create_engine来建立数据库连接写法更通用。pymysql: MySQL数据库的Python驱动。如果你用PostgreSQL或SQLite则对应安装psycopg2或sqlite3后者通常内置。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;q0.9,en;q0.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, paramsparams, headersheaders, timeout10) 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(columnscolumn_mapping, inplaceTrue) # 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’, ascendingFalse, inplaceTrue) # 按日期倒序最新在前 df.reset_index(dropTrue, inplaceTrue) # 6. 添加数据获取时间戳 df[‘created_at’] datetime.now() return df实操心得在数据清洗时一定要先打印出几行原始数据看看格式。东方财富的数值字段有时会是以“万”或“亿”为单位的字符串或者带有百分号。pd.to_numeric的errors‘coerce’参数非常有用它会把转换失败的值变成NaN而不是让整个程序崩溃方便后续排查。4. 数据库设计为股票资金流向数据安家数据抓取和清洗完成后下一步就是持久化存储。选择MySQL作为数据库主要是因为它普及率高、生态成熟与Python通过pymysqlsqlalchemy和很多数据分析工具如pandas、Metabase集成起来非常方便。4.1 表结构设计思路设计表结构时要考虑以下几点唯一性约束同一个股票在同一天的数据不应该重复存储。最自然的唯一键是(symbol, trade_date)。字段类型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。索引为了加快按股票代码和日期范围的查询速度必须在(symbol, trade_date)上建立复合索引并且由于它是唯一键本身就带有索引。如果经常需要按单个字段查询可以考虑单独为trade_date建立索引。元信息添加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) -- 按日期查询的索引 ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 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’]) # 构建连接字符串 # 格式 dialectdriver://username:passwordhost:port/database connection_string f“mysqlpymysql://{db_config[‘user’]}:{encoded_password}{db_config[‘host’]}:{db_config[‘port’]}/{db_config[‘database’]}” # 创建引擎echoTrue可以在控制台看到执行的SQL调试时有用 engine create_engine(connection_string, echoFalse) return engine5. 数据入库策略增量更新与避免重复最直接的入库方式是把抓取到的所有数据一次性插入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, conengine, if_exists‘replace’, indexFalse) # 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 错误处理与重试机制网络请求和数据库操作都可能失败。必须添加健壮的错误处理和重试逻辑。对于网络请求可以使用tenacity或retrying库实现指数退避重试。from tenacity import retry, stop_after_attempt, wait_exponential retry(stopstop_after_attempt(3), waitwait_exponential(multiplier1, min2, max10)) 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.ini或config.yaml或环境变量来管理。6.4 定时任务与自动化要让数据每天自动更新可以结合操作系统的定时任务如Linux的cron或Windows的任务计划程序来定期执行这个Python脚本。更复杂的调度可以使用Airflow、Prefect这样的工作流管理工具。6.5 数据验证与监控数据入库后可以添加简单的验证步骤比如检查最新日期的数据是否已成功入库或者统计每日入库的数据量是否在合理范围内。这可以通过在脚本最后执行一些查询语句来实现。6.6 扩展性考虑更多数据源东方财富的接口可能不稳定或有限制。可以同时集成其他免费数据源如AKShare库它封装了多个财经数据接口作为备份或补充。更多数据维度除了资金流向还可以并行抓取日K线、财务指标、龙虎榜等数据丰富你的数据库。数据质量定期检查数据中的异常值如涨跌幅超过限制、资金流数据为0的交易日等并建立清洗规则。这个项目从简单的需求出发串联起了网络爬虫、数据清洗、数据库操作、错误处理等多个核心技能点。把它跑通并不断优化你收获的不仅仅是一堆股票数据更是一套处理数据管道问题的实战经验。在实际操作中最大的挑战往往不是代码本身而是目标网站的反爬策略变化这就需要你保持对网络请求的敏锐观察随时准备调整你的爬虫策略。