195 lines
6.7 KiB
Python
195 lines
6.7 KiB
Python
#!/usr/bin/env python3
|
||
"""
|
||
统一数据获取模块 - 自动降级数据源
|
||
当东方财富API被封(腾讯云等环境)时,自动切换腾讯财经/新浪等备用数据源
|
||
|
||
使用方法(替代 akshare 直接调用):
|
||
from utils.data_fetcher import fetch_stock_hist
|
||
df = fetch_stock_hist('000001', period='daily', start_date='20250101', end_date='20260226', adjust='qfq')
|
||
"""
|
||
import logging
|
||
import requests
|
||
import pandas as pd
|
||
from datetime import datetime, timedelta
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# 数据源状态追踪 (避免反复尝试已知失败的数据源)
|
||
_source_status = {
|
||
'eastmoney': True, # 是否可用
|
||
'tencent': True,
|
||
'sina': True,
|
||
}
|
||
_source_fail_count = {
|
||
'eastmoney': 0,
|
||
'tencent': 0,
|
||
'sina': 0,
|
||
}
|
||
_MAX_FAIL_BEFORE_SKIP = 3 # 连续失败N次后暂时跳过
|
||
|
||
|
||
def _mark_source_failed(source):
|
||
"""标记数据源失败"""
|
||
_source_fail_count[source] = _source_fail_count.get(source, 0) + 1
|
||
if _source_fail_count[source] >= _MAX_FAIL_BEFORE_SKIP:
|
||
_source_status[source] = False
|
||
logger.warning(f"[数据源] {source} 连续失败 {_source_fail_count[source]} 次,暂时禁用")
|
||
|
||
|
||
def _mark_source_ok(source):
|
||
"""标记数据源成功"""
|
||
_source_fail_count[source] = 0
|
||
_source_status[source] = True
|
||
|
||
|
||
def _to_tencent_symbol(stock_code):
|
||
"""转为腾讯API格式: sz000001, sh600519"""
|
||
code = str(stock_code).strip()
|
||
if code.startswith('6'):
|
||
return f'sh{code}'
|
||
elif code.startswith('0') or code.startswith('3'):
|
||
return f'sz{code}'
|
||
elif code.startswith('8') or code.startswith('4'):
|
||
return f'bj{code}'
|
||
return f'sz{code}'
|
||
|
||
|
||
def _fetch_hist_from_tencent(stock_code, start_date, end_date, adjust='qfq'):
|
||
"""
|
||
从腾讯财经获取历史日K线
|
||
API: http://web.ifzq.gtimg.cn/appstock/app/fqkline/get
|
||
返回格式: [date, open, close, high, low, volume]
|
||
注意: 腾讯返回的是 [open, close],akshare返回的是 [开盘, 收盘]
|
||
"""
|
||
symbol = _to_tencent_symbol(stock_code)
|
||
|
||
# 腾讯最多返回约640条日线(约2.5年)
|
||
# 格式化日期
|
||
start_fmt = f'{start_date[:4]}-{start_date[4:6]}-{start_date[6:8]}' if len(start_date) == 8 else start_date
|
||
end_fmt = f'{end_date[:4]}-{end_date[4:6]}-{end_date[6:8]}' if len(end_date) == 8 else end_date
|
||
|
||
# 计算请求的天数
|
||
try:
|
||
d1 = datetime.strptime(start_date[:8], '%Y%m%d')
|
||
d2 = datetime.strptime(end_date[:8], '%Y%m%d')
|
||
num_bars = (d2 - d1).days + 50 # 多请求一些,因为有非交易日
|
||
num_bars = min(max(num_bars, 60), 640)
|
||
except:
|
||
num_bars = 320
|
||
|
||
# 前复权: qfqday, 不复权: day
|
||
adj_key = 'qfqday' if adjust == 'qfq' else 'day'
|
||
adj_param = 'qfq' if adjust == 'qfq' else ''
|
||
|
||
url = f'http://web.ifzq.gtimg.cn/appstock/app/fqkline/get'
|
||
params = f'{symbol},day,{start_fmt},{end_fmt},{num_bars},{adj_param}'
|
||
|
||
r = requests.get(url, params={'param': params}, timeout=15)
|
||
if r.status_code != 200:
|
||
raise Exception(f'腾讯API返回 {r.status_code}')
|
||
|
||
data = r.json()
|
||
stock_key = symbol # e.g. 'sz000001'
|
||
klines = data.get('data', {}).get(stock_key, {}).get(adj_key, [])
|
||
|
||
if not klines:
|
||
# 尝试不复权
|
||
klines = data.get('data', {}).get(stock_key, {}).get('day', [])
|
||
|
||
if not klines:
|
||
return pd.DataFrame()
|
||
|
||
# 构造与 akshare stock_zh_a_hist 兼容的 DataFrame
|
||
# 腾讯格式: [date, open, close, high, low, volume]
|
||
rows = []
|
||
prev_close = None
|
||
for k in klines:
|
||
if len(k) < 6:
|
||
continue
|
||
date_str = k[0]
|
||
open_price = float(k[1])
|
||
close_price = float(k[2])
|
||
high_price = float(k[3])
|
||
low_price = float(k[4])
|
||
volume = float(k[5])
|
||
|
||
# 计算衍生字段
|
||
change_amount = close_price - prev_close if prev_close else 0
|
||
change_pct = (change_amount / prev_close * 100) if prev_close and prev_close > 0 else 0
|
||
amplitude = ((high_price - low_price) / prev_close * 100) if prev_close and prev_close > 0 else 0
|
||
|
||
rows.append({
|
||
'日期': date_str,
|
||
'开盘': open_price,
|
||
'收盘': close_price,
|
||
'最高': high_price,
|
||
'最低': low_price,
|
||
'成交量': int(volume),
|
||
'成交额': 0, # 腾讯不提供成交额
|
||
'振幅': round(amplitude, 2),
|
||
'涨跌幅': round(change_pct, 2),
|
||
'涨跌额': round(change_amount, 2),
|
||
'换手率': 0, # 腾讯不提供换手率
|
||
})
|
||
prev_close = close_price
|
||
|
||
df = pd.DataFrame(rows)
|
||
|
||
# 过滤日期范围
|
||
if not df.empty:
|
||
df['日期'] = pd.to_datetime(df['日期'])
|
||
start_dt = pd.to_datetime(start_fmt)
|
||
end_dt = pd.to_datetime(end_fmt)
|
||
df = df[(df['日期'] >= start_dt) & (df['日期'] <= end_dt)]
|
||
df = df.sort_values('日期').reset_index(drop=True)
|
||
|
||
return df
|
||
|
||
|
||
def fetch_stock_hist(stock_code, period='daily', start_date='20200101',
|
||
end_date=None, adjust='qfq'):
|
||
"""
|
||
获取股票历史K线数据(腾讯财经为主数据源)
|
||
|
||
参数与 akshare.stock_zh_a_hist 完全兼容:
|
||
stock_code: 股票代码 (纯数字,如 '000001')
|
||
period: 'daily', 'weekly', 'monthly'
|
||
start_date: 开始日期 'YYYYMMDD'
|
||
end_date: 结束日期 'YYYYMMDD'
|
||
adjust: 'qfq'(前复权) / 'hfq'(后复权) / ''(不复权)
|
||
|
||
返回: pandas DataFrame, 与 akshare 格式兼容
|
||
"""
|
||
if end_date is None:
|
||
end_date = datetime.now().strftime('%Y%m%d')
|
||
|
||
# 数据源: 腾讯财经(日K线,腾讯云最快最稳)
|
||
if period == 'daily':
|
||
try:
|
||
df = _fetch_hist_from_tencent(stock_code, start_date, end_date, adjust)
|
||
if df is not None and not df.empty:
|
||
return df
|
||
except Exception as e:
|
||
logger.warning(f"[数据源] 腾讯财经获取失败({stock_code}): {str(e)[:100]}")
|
||
|
||
logger.error(f"[数据源] 数据源获取失败: {stock_code}")
|
||
return pd.DataFrame()
|
||
|
||
|
||
def fetch_stock_codes():
|
||
"""
|
||
获取全部A股股票代码和名称(从数据库获取)
|
||
返回: DataFrame with columns ['code', 'name']
|
||
"""
|
||
# 从数据库获取(最可靠,不依赖外部API)
|
||
logger.info("[数据源] 从数据库获取股票列表")
|
||
return pd.DataFrame()
|
||
|
||
|
||
def reset_source_status():
|
||
"""重置所有数据源状态(用于定时任务开始时)"""
|
||
global _source_status, _source_fail_count
|
||
_source_status = {'eastmoney': True, 'tencent': True, 'sina': True}
|
||
_source_fail_count = {'eastmoney': 0, 'tencent': 0, 'sina': 0}
|
||
logger.info("[数据源] 所有数据源状态已重置")
|