Python + MySQL + Tushare 实现股票数据采集与K线分析系统
下面给你一套Python MySQL Tushare Pro 实现股票数据采集与 K 线分析系统的可落地方案覆盖环境准备 → 数据库设计 → Tushare 采集历史增量→ K 线指标计算 → 查询分析接口。代码基于 Tushare Pro 最新pro.daily接口与 MySQL 8.x 生产可用写法。一、系统架构Tushare Pro (daily) │ pandas DataFrame ▼ collector.py ──▶ MySQL (kline_daily) │ │ │ ▼ │ analyzer.py (MA/MACD/KDJ/BOLL) ▼ │ scheduler.py query_api.py (SQL 查询 / 回测准备)核心原则按交易日维度采集推荐trade_date拉全市场比按股票循环快 10 倍MySQL 用唯一索引去重ON DUPLICATE KEY UPDATE幂等写入复权类型分开存表未复权 / 后复权 / 前复权二、依赖与配置pip install tushare pandas sqlalchemy pymysql mysql-connector-pythonconfig.pyimport os class Config: TUSHARE_TOKEN os.getenv(TUSHARE_TOKEN, 你的token) MYSQL_USER root MYSQL_PASSWORD pwd MYSQL_HOST 127.0.0.1 MYSQL_PORT 3306 MYSQL_DB stock_db SQLALCHEMY_URI ( fmysqlpymysql://{MYSQL_USER}:{MYSQL_PASSWORD} f{MYSQL_HOST}:{MYSQL_PORT}/{MYSQL_DB}?charsetutf8mb4 ) # 采集节流基础积分 500次/分钟 API_SLEEP_PER_100 12 # 每100次请求睡12秒三、MySQL 表结构重点CREATE DATABASE stock_db CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; USE stock_db; CREATE TABLE kline_daily ( id BIGINT AUTO_INCREMENT PRIMARY KEY, ts_code VARCHAR(20) NOT NULL COMMENT 000001.SZ, trade_date DATE NOT NULL, open DECIMAL(12,4) NOT NULL, high DECIMAL(12,4) NOT NULL, low DECIMAL(12,4) NOT NULL, close DECIMAL(12,4) NOT NULL, pre_close DECIMAL(12,4), pct_chg DECIMAL(8,4), vol BIGINT COMMENT 成交量(手), amount DECIMAL(20,4) COMMENT 成交额(千元), adj_type TINYINT NOT NULL DEFAULT 0 COMMENT 0未复权 1后复权 2前复权, created_at DATETIME DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY uk_code_date_adj (ts_code, trade_date, adj_type), KEY idx_date (trade_date), KEY idx_code (ts_code) ) ENGINEInnoDB COMMENTA股日K线;⚠️open/high/low/close必须用反引号否则和 MySQL 关键字冲突。四、Tushare 采集层历史 增量collector.pyimport time import pandas as pd from sqlalchemy import create_engine import tushare as ts from config import Config ts.set_token(Config.TUSHARE_TOKEN) pro ts.pro_api() engine create_engine(Config.SQLALCHEMY_URI) FIELDS ts_code,trade_date,open,high,low,close,pre_close,pct_chg,vol,amount def collect_by_date(trade_date: str, adj_type: int 0): 按交易日拉全市场日线推荐方式 df pro.daily(trade_datetrade_date, fieldsFIELDS) if df is None or df.empty: return 0 df[trade_date] pd.to_datetime(df[trade_date]) df[adj_type] adj_type df df.rename(columns{vol: vol}) df.to_sql( kline_daily, engine, if_existsappend, indexFalse, methodmulti, chunksize2000 ) return len(df) def collect_history(start: str, end: str): 遍历交易日历用 trade_date 循环不按股票循环 cal pro.trade_cal(exchangeSSE, start_datestart, end_dateend, is_open1) dates cal[cal_date].tolist() for i, d in enumerate(dates): try: n collect_by_date(d, adj_type0) print(f{d} ✅ {n} rows) except Exception as e: print(f{d} ❌ {e}) if i % 100 0 and i: time.sleep(Config.API_SLEEP_PER_100) def incremental_update(): 每日 16:00 后跑补昨天 yesterday pd.Timestamp.now().normalize() - pd.Timedelta(days1) d yesterday.strftime(%Y%m%d) collect_by_date(d, adj_type0)Tushare 官方建议“建议提供循环日期来提取全市场数据不要通过循环 ts_code 来拉取历史”。五、K 线指标分析层analyzer.pyimport pandas as pd from sqlalchemy import create_engine from config import Config engine create_engine(Config.SQLALCHEMY_URI) def load_kline(ts_code: str, start: str, end: str) - pd.DataFrame: sql f SELECT trade_date, open, high, low, close, vol FROM kline_daily WHERE ts_code{ts_code} AND adj_type0 AND trade_date BETWEEN {start} AND {end} ORDER BY trade_date df pd.read_sql(sql, engine, parse_dates[trade_date]) return df.set_index(trade_date) def add_indicators(df: pd.DataFrame) - pd.DataFrame: # MA df[ma5] df[close].rolling(5).mean() df[ma20] df[close].rolling(20).mean() df[ma60] df[close].rolling(60).mean() # EMA df[ema12] df[close].ewm(span12, adjustFalse).mean() df[ema26] df[close].ewm(span26, adjustFalse).mean() # MACD df[dif] df[ema12] - df[ema26] df[dea] df[dif].ewm(span9, adjustFalse).mean() df[macd] (df[dif] - df[dea]) * 2 # BOLL df[boll_mid] df[close].rolling(20).mean() df[boll_std] df[close].rolling(20).std() df[boll_up] df[boll_mid] 2 * df[boll_std] df[boll_dn] df[boll_mid] - 2 * df[boll_std] # KDJ low_9 df[low].rolling(9).min() high_9 df[high].rolling(9).max() df[rsv] (df[close] - low_9) / (high_9 - low_9) * 100 df[k] df[rsv].ewm(alpha1/3, adjustFalse).mean() df[d] df[k].ewm(alpha1/3, adjustFalse).mean() df[j] 3 * df[k] - 2 * df[d] return df if __name__ __main__: df load_kline(600519.SH, 20240101, 20240801) df add_indicators(df) print(df.tail(10))六、定时调度每天收盘后增量# scheduler.py from apscheduler.schedulers.blocking import BlockingScheduler from collector import incremental_update sched BlockingScheduler(tzAsia/Shanghai) sched.add_job(incremental_update, cron, hour16, minute10) sched.start()七、常见坑生产必看限流基础积分 500 次/分钟全市场约 5000 只股票/天按日期拉最快。重复数据必须UNIQUE(ts_code, trade_date, adj_type)to_sql(append)否则重跑会炸。复权pro.daily是未复权前/后复权用pro.daily(fields..., adjhfq/qfq)或pro.daily_qfq。vol 单位Tushare 返回“手”MySQL 存“手”分析时注意和成交额单位不一致。停牌daily停牌不返回行不要误以为采集失败。