Initial commit

This commit is contained in:
freedakgmail
2026-07-17 18:49:07 +08:00
commit 9c7d7abdd4
100 changed files with 41337 additions and 0 deletions
+715
View File
@@ -0,0 +1,715 @@
#!/usr/bin/env python3
"""
5分钟K线数据每日采集脚本
功能:每日收盘后自动采集全市场A股的5分钟K线数据,存入 stock_kline_5min 表。
数据源:akshare stock_zh_a_minute(新浪财经,免费,腾讯云可用)
特点:
- 增量采集:只采集当日新增数据
- 断点续传:记录已采集的股票,中断后可继续
- 频率控制:自动限速避免触发API封禁
- 回填模式:可手动回填最近N天的历史分钟数据
用法:
# 每日采集(推荐在 17:00 后运行,收盘后数据完整)
./venv/bin/python sync_kline_5min.py
# 回填最近5天
./venv/bin/python sync_kline_5min.py --backfill 5
# 只采集指定股票
./venv/bin/python sync_kline_5min.py --codes 300720,000001
# 快速测试(只采集前10只)
./venv/bin/python sync_kline_5min.py --limit 10
建议定时任务:
30 17 * * 1-5 /opt/stock-app/venv/bin/python /opt/stock-app/sync_kline_5min.py >> /opt/stock-app/sync_kline_5min.log 2>&1
"""
import sys
import os
import time
import argparse
import signal as sig_module
import fcntl
import atexit
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, date, timedelta
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import psycopg2
from psycopg2.extras import execute_values
from config import Config
# ============ 配置 ============
API_DELAY = 0.8 # 每次API调用间隔(秒),东财限流严格需更保守
BATCH_SAVE_SIZE = 500 # 每批保存行数
PERIOD = '5' # K线级别:'1','5','15','30','60'
MAX_RETRIES = 2 # 单只股票API重试次数
RETRY_DELAY = 3 # 重试等待时间(秒)
EM_CIRCUIT_BREAKER = 3 # 东财API连续失败N次后暂停使用(快速熔断)
LOCK_FILE = '/tmp/sync_kline_5min.lock'
DEFAULT_WORKERS = 3 # 默认并发数(东财API限流严格,不宜过高)
_shutdown = False
_em_consecutive_errors = 0 # 东财API连续错误计数
_em_disabled = False # 东财API是否被暂停
_em_lock = threading.Lock() # 东财API状态锁
_lock_fd = None # 进程锁文件描述符
def _clean_stale_lock():
"""清理残留的锁文件(进程已不存在或权限不对时)"""
if not os.path.exists(LOCK_FILE):
return
try:
with open(LOCK_FILE, 'r') as f:
old_pid = f.read().strip()
if old_pid and old_pid.isdigit():
try:
os.kill(int(old_pid), 0) # 检查进程是否存在
return # 进程仍在运行,不清理
except ProcessLookupError:
pass # 进程已不存在,清理
except PermissionError:
return # 进程存在但无权限检查,不清理
except (IOError, PermissionError):
pass # 无法读取锁文件,尝试删除
try:
os.remove(LOCK_FILE)
print(f"🧹 已清理残留锁文件 (旧PID: {old_pid if 'old_pid' in dir() else '未知'})", flush=True)
except Exception:
pass
def acquire_lock():
"""获取进程锁,防止多实例同时运行"""
global _lock_fd
# 先尝试清理残留的锁文件
_clean_stale_lock()
try:
_lock_fd = open(LOCK_FILE, 'w')
except PermissionError:
# 锁文件权限不对,尝试删除后重建
try:
os.remove(LOCK_FILE)
_lock_fd = open(LOCK_FILE, 'w')
except Exception as e:
print(f"⚠️ 无法创建锁文件 {LOCK_FILE}: {e}", flush=True)
return False
try:
fcntl.flock(_lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
_lock_fd.write(str(os.getpid()))
_lock_fd.flush()
atexit.register(release_lock)
return True
except IOError:
# 另一个实例正在运行,读取其PID
try:
with open(LOCK_FILE, 'r') as f:
old_pid = f.read().strip()
print(f"⚠️ 另一个实例正在运行 (PID: {old_pid}),退出", flush=True)
except Exception:
print(f"⚠️ 另一个实例正在运行,退出", flush=True)
_lock_fd.close()
_lock_fd = None
return False
def release_lock():
"""释放进程锁"""
global _lock_fd
if _lock_fd:
try:
fcntl.flock(_lock_fd, fcntl.LOCK_UN)
_lock_fd.close()
except Exception:
pass
_lock_fd = None
try:
os.remove(LOCK_FILE)
except Exception:
pass
def signal_handler(signum, frame):
global _shutdown
print("\n⚠️ 收到中断信号,正在优雅退出...", flush=True)
_shutdown = True
sig_module.signal(sig_module.SIGINT, signal_handler)
sig_module.signal(sig_module.SIGTERM, signal_handler)
def get_db_conn():
return psycopg2.connect(
host=Config.DB_HOST, port=Config.DB_PORT,
dbname=Config.DB_NAME, user=Config.DB_USER, password=Config.DB_PASSWORD,
)
def ensure_table(conn):
"""确保 stock_kline_5min 表存在"""
with conn.cursor() as cur:
cur.execute("""
CREATE TABLE IF NOT EXISTS stock_kline_5min (
code VARCHAR(10) NOT NULL,
dt TIMESTAMP NOT NULL,
open DECIMAL(12, 4),
high DECIMAL(12, 4),
low DECIMAL(12, 4),
close DECIMAL(12, 4),
volume BIGINT,
amount DECIMAL(20, 2),
change_pct DECIMAL(8, 4),
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (code, dt)
)
""")
cur.execute("CREATE INDEX IF NOT EXISTS idx_kline_5min_dt ON stock_kline_5min(dt)")
conn.commit()
def get_stock_codes(conn, only_codes=None):
"""获取需要采集的股票列表"""
with conn.cursor() as cur:
if only_codes:
placeholders = ','.join(['%s'] * len(only_codes))
cur.execute(f"SELECT code, name FROM stock_realtime_price WHERE code IN ({placeholders}) ORDER BY code",
only_codes)
else:
# 获取所有有效股票(价格>0的)
cur.execute("""
SELECT code, name FROM stock_realtime_price
WHERE price > 0 AND code NOT LIKE 'BJ%%'
ORDER BY code
""")
return cur.fetchall()
def get_already_synced_codes(conn, target_date):
"""获取今天已经同步过的股票(用于断点续传)"""
with conn.cursor() as cur:
cur.execute("""
SELECT DISTINCT code FROM stock_kline_5min
WHERE dt::date = %s
""", (target_date,))
return {r[0] for r in cur.fetchall()}
def _code_to_sina_symbol(code):
"""股票代码转新浪格式:000001 → sz000001, 600519 → sh600519"""
if code.startswith(('0', '3')):
return f'sz{code}'
elif code.startswith(('6', '5')):
return f'sh{code}'
elif code.startswith(('8', '9', '4')):
return f'bj{code}'
return f'sz{code}'
def fetch_5min_kline_sina(code):
"""
数据源1:新浪APIakshare stock_zh_a_minute
优点:稳定、不易被封、回溯约2个月
返回: list of tuple (code, dt, open, high, low, close, volume, amount, change_pct)
返回 None 表示无数据(非错误)
返回 'error' 字符串表示API错误
"""
import akshare as ak
for attempt in range(MAX_RETRIES + 1):
try:
symbol = _code_to_sina_symbol(code)
df = ak.stock_zh_a_minute(symbol=symbol, period=PERIOD)
if df is None or df.empty:
return None # 无数据,非错误
rows = []
for _, row in df.iterrows():
dt_str = str(row.get('day', ''))
try:
dt = datetime.strptime(dt_str, '%Y-%m-%d %H:%M:%S')
except ValueError:
continue
o = float(row.get('open', 0))
h = float(row.get('high', 0))
l = float(row.get('low', 0))
c = float(row.get('close', 0))
v = int(float(row.get('volume', 0)))
rows.append((code, dt, o, h, l, c, v, 0.0, 0.0))
return rows if rows else None
except (IndexError, KeyError, ValueError):
# list index out of range / KeyError / ValueError
# 这些是"该股票无5分钟数据"的表现,不是API故障
return None # 无数据,非错误
except Exception as e:
err_str = str(e)
if attempt < MAX_RETRIES:
wait = RETRY_DELAY * (attempt + 1)
print(f" ⚠️ {code}(新浪) 第{attempt+1}次失败: {err_str[:60]}, {wait}s后重试", flush=True)
time.sleep(wait)
continue
return 'error' # 真正的API错误
return 'error'
def fetch_5min_kline_em(code, start_date=None, end_date=None):
"""
数据源2:东方财富APIakshare stock_zh_a_hist_min_em
优点:有成交额和涨跌幅,回溯约2个月
缺点:容易被限流
返回: list of tuple 或 None(无数据) 或 'error'API错误)
"""
import akshare as ak
for attempt in range(MAX_RETRIES + 1):
try:
kwargs = {'symbol': code, 'period': PERIOD, 'adjust': ''}
if start_date:
kwargs['start_date'] = start_date
if end_date:
kwargs['end_date'] = end_date
df = ak.stock_zh_a_hist_min_em(**kwargs)
if df is None or df.empty:
return None
rows = []
for _, row in df.iterrows():
dt_str = str(row['时间'])
try:
dt = datetime.strptime(dt_str, '%Y-%m-%d %H:%M:%S')
except ValueError:
continue
rows.append((
code, dt,
float(row.get('开盘', 0)),
float(row.get('最高', 0)),
float(row.get('最低', 0)),
float(row.get('收盘', 0)),
int(float(row.get('成交量', 0))),
float(row.get('成交额', 0)),
float(row.get('涨跌幅', 0)),
))
return rows if rows else None
except (IndexError, KeyError, ValueError):
return None # 无数据,非错误
except Exception as e:
err_str = str(e)
# 连接被断开、限流等属于真正的API错误
if attempt < MAX_RETRIES:
wait = RETRY_DELAY * (attempt + 1) * 2 # 东财限流严重,加长等待
print(f" ⚠️ {code}(东财) 第{attempt+1}次失败: {err_str[:60]}, {wait}s后重试", flush=True)
time.sleep(wait)
continue
return 'error'
return 'error'
def fetch_5min_kline_tencent(code):
"""
数据源3:腾讯财经1分钟数据 → 聚合为5分钟K线
优点:腾讯云服务器永不被封
缺点:只有当天数据
返回: list of tuple 或 None 或 'error'
"""
import requests as _requests
import json as _json
try:
symbol = _code_to_sina_symbol(code) # sh/sz 格式通用
url = f"https://web.ifzq.gtimg.cn/appstock/app/minute/query?code={symbol}"
r = _requests.get(url, timeout=15, headers={
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
"Referer": "https://stockapp.finance.qq.com",
})
if r.status_code != 200:
return 'error'
text = r.text
start_idx = text.find("=") + 1
end_idx = text.rfind("}") + 1
data = _json.loads(text[start_idx:end_idx])
records = data.get("data", {}).get(symbol, {}).get("data", {}).get("data", [])
if not records:
return None
today = date.today()
# 解析1分钟数据: "0930 1466.99 153 22444946.66"
# 格式: HHMM price volume amount(累积)
min_data = []
for rec in records:
parts = rec.split()
if len(parts) < 4:
continue
hhmm = parts[0]
price = float(parts[1])
vol = int(parts[2])
try:
h, m = int(hhmm[:2]), int(hhmm[2:])
dt = datetime(today.year, today.month, today.day, h, m, 0)
except (ValueError, IndexError):
continue
min_data.append((dt, price, vol))
if not min_data:
return None
# 聚合为5分钟K线
# 5分钟窗口: 09:30-09:35, 09:35-09:40, ...
from collections import defaultdict
bars = defaultdict(list)
for dt, price, vol in min_data:
# 5分钟窗口起始时间
minute = dt.minute
bar_min = (minute // 5) * 5
bar_dt = dt.replace(minute=bar_min, second=0)
bars[bar_dt].append((price, vol))
rows = []
prev_vol = 0
for bar_dt in sorted(bars.keys()):
ticks = bars[bar_dt]
o = ticks[0][0] # 第一个价格
c = ticks[-1][0] # 最后一个价格
h = max(p for p, _ in ticks)
l = min(p for p, _ in ticks)
# 腾讯的volume是累积值,取窗口最后的 - 窗口最前的之前
last_vol = ticks[-1][1]
bar_vol = last_vol - prev_vol if prev_vol > 0 else ticks[-1][1]
prev_vol = last_vol
rows.append((code, bar_dt, o, h, l, c, max(0, bar_vol), 0.0, 0.0))
return rows if rows else None
except (IndexError, KeyError, ValueError):
return None
except Exception:
return 'error'
def fetch_5min_kline(code, start_date=None, end_date=None):
"""
主入口:优先东财API → 新浪API → 腾讯聚合(仅当天)
返回:
- list of tuple: 成功获取数据
- None: 该股票无5分钟数据(停牌、退市等,非错误)
- 'error': API故障/限流
"""
global _em_consecutive_errors, _em_disabled
# 1) 优先使用东财API(数据最全)
with _em_lock:
em_ok = not _em_disabled
if em_ok:
result = fetch_5min_kline_em(code, start_date, end_date)
if result == 'error':
with _em_lock:
_em_consecutive_errors += 1
if _em_consecutive_errors >= EM_CIRCUIT_BREAKER:
_em_disabled = True
print(f" ⚠️ 东财API连续失败{EM_CIRCUIT_BREAKER}次,尝试备用源", flush=True)
time.sleep(2)
else:
with _em_lock:
_em_consecutive_errors = 0
return result
# 2) 备用:新浪API
result = fetch_5min_kline_sina(code)
if result is not None and result != 'error':
return result
# 3) 终极备用:腾讯1分钟聚合(仅当天数据,但永不被封)
return fetch_5min_kline_tencent(code)
def save_batch(conn, all_rows):
"""批量保存5分钟K线数据(UPSERT"""
if not all_rows:
return 0
with conn.cursor() as cur:
execute_values(
cur,
"""
INSERT INTO stock_kline_5min (code, dt, open, high, low, close, volume, amount, change_pct)
VALUES %s
ON CONFLICT (code, dt) DO UPDATE SET
open = EXCLUDED.open,
high = EXCLUDED.high,
low = EXCLUDED.low,
close = EXCLUDED.close,
volume = EXCLUDED.volume,
amount = EXCLUDED.amount,
change_pct = EXCLUDED.change_pct,
updated_at = CURRENT_TIMESTAMP
""",
all_rows,
page_size=1000,
)
conn.commit()
return len(all_rows)
def filter_rows_by_date(rows, target_date):
"""只保留目标日期的行"""
if not rows:
return None
filtered = [r for r in rows if r[1].date() == target_date]
return filtered if filtered else None
def _probe_apis():
"""启动时探测各API是否可用,提前设置熔断状态"""
global _em_disabled, _em_consecutive_errors
print("🔍 探测数据源可用性...", flush=True)
# 测试东财API(用最活跃的股票)
em_ok = False
try:
import akshare as ak
df = ak.stock_zh_a_hist_min_em(symbol='600519', period='5', adjust='')
if df is not None and not df.empty:
em_ok = True
print(" ✅ 东财API: 可用", flush=True)
else:
print(" ❌ 东财API: 返回空数据", flush=True)
except Exception as e:
print(f" ❌ 东财API: {str(e)[:60]}", flush=True)
if not em_ok:
with _em_lock:
_em_disabled = True
_em_consecutive_errors = EM_CIRCUIT_BREAKER
print(" → 东财API已禁用,将使用备用源", flush=True)
# 测试新浪API
sina_ok = False
try:
result = fetch_5min_kline_sina('600519')
if result is not None and result != 'error':
sina_ok = True
print(" ✅ 新浪API: 可用", flush=True)
else:
print(" ❌ 新浪API: 不可用", flush=True)
except Exception:
print(" ❌ 新浪API: 异常", flush=True)
# 腾讯API(聚合方式)总是可用
print(" ✅ 腾讯API: 始终可用(聚合1分钟→5分钟)", flush=True)
source = "东财" if em_ok else ("新浪" if sina_ok else "腾讯(聚合)")
print(f" 📡 主数据源: {source}", flush=True)
return em_ok, sina_ok
def main():
global _shutdown
# 进程锁 — 防止多实例同时运行
if not acquire_lock():
sys.exit(1)
parser = argparse.ArgumentParser(description='5分钟K线数据采集')
parser.add_argument('--backfill', type=int, default=0, metavar='DAYS',
help='回填最近N天的数据(默认0=只采集当天)')
parser.add_argument('--codes', type=str, default=None,
help='只采集指定股票(逗号分隔,如 300720,000001')
parser.add_argument('--limit', type=int, default=0,
help='限制采集股票数量(0=全部,用于测试)')
parser.add_argument('--delay', type=float, default=API_DELAY,
help=f'API调用间隔秒数(默认{API_DELAY}')
parser.add_argument('--resume', action='store_true',
help='跳过今天已采集的股票(断点续传)')
parser.add_argument('--workers', type=int, default=DEFAULT_WORKERS,
help=f'并发线程数(默认{DEFAULT_WORKERS}')
args = parser.parse_args()
today = date.today()
target_dates = [today]
if args.backfill > 0:
for i in range(1, args.backfill + 1):
d = today - timedelta(days=i)
if d.weekday() < 5: # 跳过周末
target_dates.append(d)
target_dates.sort()
only_codes = args.codes.split(',') if args.codes else None
print(f"{'='*70}", flush=True)
print(f"📊 5分钟K线数据采集(并发模式)", flush=True)
print(f"📅 目标日期: {', '.join(str(d) for d in target_dates)}", flush=True)
print(f"⏱️ API间隔: {args.delay}s | 并发: {args.workers} 线程", flush=True)
if only_codes:
print(f"🎯 指定股票: {only_codes}", flush=True)
print(f"{'='*70}", flush=True)
# 启动时探测API可用性,提前熔断不可用的源
em_ok, sina_ok = _probe_apis()
# 如果主源不可用,腾讯聚合模式可以用更多并发(不限流)
if not em_ok and not sina_ok:
if args.workers < 5:
args.workers = 5
print(f" 📡 腾讯模式:提升并发到 {args.workers} 线程", flush=True)
if args.delay > 0.3:
args.delay = 0.3
print(f" 📡 腾讯模式:降低延迟到 {args.delay}s", flush=True)
conn = get_db_conn()
ensure_table(conn)
all_stocks = get_stock_codes(conn, only_codes)
if args.limit > 0:
all_stocks = all_stocks[:args.limit]
total = len(all_stocks)
print(f"📈 待采集股票: {total}", flush=True)
if args.resume:
synced = get_already_synced_codes(conn, today)
before = len(all_stocks)
all_stocks = [(c, n) for c, n in all_stocks if c not in synced]
print(f"🔄 断点续传: 跳过 {before - len(all_stocks)} 只已同步, 剩余 {len(all_stocks)}", flush=True)
total = len(all_stocks)
start_time = time.time()
done_count = 0
success_count = 0
error_count = 0
skip_count = 0
total_rows = 0
save_buffer = []
last_report_time = time.time()
# 线程安全锁
_stats_lock = threading.Lock()
_buffer_lock = threading.Lock()
def _fetch_one(code_name):
"""单只股票采集任务(在工作线程中运行)"""
code, name = code_name
if _shutdown:
return None
# 线程内延迟,分散API请求
time.sleep(args.delay)
result = fetch_5min_kline(code)
return (code, name, result)
print(f"\n🚀 开始采集({args.workers}线程并发)...", flush=True)
print(f"-" * 70, flush=True)
with ThreadPoolExecutor(max_workers=args.workers) as executor:
futures = {executor.submit(_fetch_one, item): item for item in all_stocks}
for future in as_completed(futures):
if _shutdown:
break
ret = future.result()
if ret is None:
continue
code, name, result = ret
with _stats_lock:
done_count += 1
if isinstance(result, list):
rows = result
if args.backfill == 0:
rows = filter_rows_by_date(rows, today)
if rows:
with _buffer_lock:
save_buffer.extend(rows)
success_count += 1
else:
skip_count += 1
elif result == 'error':
error_count += 1
else:
skip_count += 1
# 攒够一批就保存
with _buffer_lock:
if len(save_buffer) >= BATCH_SAVE_SIZE:
saved = save_batch(conn, save_buffer)
total_rows += saved
save_buffer = []
# 进度报告
now = time.time()
if now - last_report_time >= 5:
elapsed = now - start_time
speed = done_count / elapsed if elapsed > 0 else 0
remaining = (total - done_count) / speed if speed > 0 else 0
pct = done_count / total * 100 if total > 0 else 100
print(f" [{pct:5.1f}%] {done_count}/{total} "
f"| {speed:.1f}只/秒 | 剩余 {remaining/60:.1f}分钟 "
f"| ✅{success_count}{error_count} ⏭️{skip_count} "
f"| 已保存 {total_rows}", flush=True)
last_report_time = now
# 保存剩余数据
if save_buffer:
saved = save_batch(conn, save_buffer)
total_rows += saved
elapsed = time.time() - start_time
speed = done_count / elapsed if elapsed > 0 else 0
print(f"\n{'='*70}", flush=True)
print(f"{'⏹️ 中断' if _shutdown else '✅ 完成'}", flush=True)
print(f" 处理: {done_count}/{total}", flush=True)
print(f" 成功: {success_count} | 失败: {error_count} | 跳过: {skip_count}", flush=True)
print(f" 保存: {total_rows} 条 5分钟K线", flush=True)
print(f" 耗时: {elapsed/60:.1f} 分钟 ({speed:.1f} 只/秒)", flush=True)
# 显示数据库统计
with conn.cursor() as cur:
cur.execute("""
SELECT count(*), count(DISTINCT code),
min(dt::date), max(dt::date),
count(DISTINCT dt::date)
FROM stock_kline_5min
""")
cnt, codes, min_d, max_d, days = cur.fetchone()
print(f"\n📊 stock_kline_5min 数据库统计:", flush=True)
print(f" 总记录: {cnt:,} 条 | 覆盖: {codes} 只股票 | {days}", flush=True)
print(f" 日期: {min_d} ~ {max_d}", flush=True)
# 按日期统计
cur.execute("""
SELECT dt::date AS trade_date, count(*), count(DISTINCT code)
FROM stock_kline_5min
GROUP BY dt::date
ORDER BY dt::date DESC
LIMIT 5
""")
print(f" 最近5天:", flush=True)
for d, cnt, codes in cur.fetchall():
print(f" {d}: {cnt:>8,} 条 ({codes} 只股票)", flush=True)
print(f"{'='*70}", flush=True)
conn.close()
if __name__ == '__main__':
main()