""" live_data.py — 全特征 API 实时更新(不依赖本地 CSV) ===================================================== 数据源: 1. FRED API — 美国宏观、价格、利率、ISM 子指数 (~90) 2. Yahoo Finance — 商品价格、汇率、波动率指数 (~25) 3. World Bank API — 全球大宗商品价格 (~40) 4. akshare — 中国宏观 M1/M2/PMI (~9) 5. CFTC — 持仓报告(自动下载 CSV)(~30) 6. EIA API — 美国能源库存/产量 (~8) 7. GPR Index — 地缘政治风险指数 (Caldara & Iacoviello) 8. 派生特征 — 从原始特征自动计算 (~55) 用法: python live_data.py # 全量更新 python live_data.py --test # 测试各 API 连通性 python live_data.py --source fred # 只更新某个数据源 """ import pandas as pd import numpy as np import os, json, warnings, argparse, time, requests from datetime import datetime, timedelta warnings.filterwarnings('ignore') os.chdir(r'e:\大三下\比赛\花旗杯\模型') FRED_KEY = 'fc02a6e6a359a4cc16f0f1752d258011' EIA_KEY = '9Nv5PhLREMmmKeo0zJ2U3Zu21Bntf8DfhEKBpi55' OUTPUT_DIR = 'output' # ═══════════════════════════════════════════════════════════ # 1. FRED API — 美国宏观 + 价格 + 利率 + ISM # ═══════════════════════════════════════════════════════════ FRED_MAP = { # ── 价格 ── 'WTI_spot': 'DCOILWTICO', 'Brent_spot': 'DCOILBRENTEU', 'natgas_spot_henry': 'DHHNGSP', 'gold_spot': 'GOLDAMGBD228NLBM', # ── 需求/宏观 ── 'pmi_us_mfg': 'MANEMP', 'ipi_us': 'INDPRO', 'nonfarm_us': 'PAYEMS', 'usd_index': 'DTWEXBGS', 'cpi_us': 'CPIAUCSL', 'fed_funds_rate': 'FEDFUNDS', 'yield_spread_10y2y': 'T10Y2Y', 'us_10y_yield': 'DGS10', # ── 风险 ── 'vix': 'VIXCLS', # ── 供给 ── 'us_oil_inventory_total': 'WCESTUS1', 'us_crude_production': 'MCRFPUS2', # ── 利率/收益率 ── 'libor_usd_3m': 'USD3MTD156N', '美国:国债收益率:3年': 'DGS3', '美国:国债收益率:7年': 'DGS7', '美国:国债收益率:10年': 'DGS10', '美国:国债收益率:30年': 'DGS30', '美国:国债到期收益率(以通胀为标的):7年': 'DFII7', '美国:国债收益率利差:10年-2年': 'T10Y2Y', # ── 货币 ── '美国:M1:非季调': 'M1NS', 'us_m2_yoy': 'WM2NS', '美国:货币乘数:M1:非季调': 'MULT', '美国:货币乘数:M2:非季调': 'M2REAL', # ── ISM 制造业 PMI 子指数(40个)── '美国:ISM:制造业PMI': 'NAPM', '美国:ISM:制造业PMI:新订单': 'NAPMNOI', '美国:ISM:制造业PMI:自有库存': 'NAPMII', '美国:ISM:制造业PMI:客户库存': 'NAPMCI', '美国:ISM:制造业PMI:就业': 'NAPMEI', '美国:ISM:制造业PMI:物价': 'NAPMPRI', '美国:ISM:制造业PMI:产出': 'NAPMPI', '美国:ISM:制造业PMI:供应商交付': 'NAPMSDI', '美国:ISM:制造业PMI:新出口订单': 'NAPMEX', '美国:ISM:制造业PMI:进口': 'NAPMIMP', '美国:ISM:服务业PMI': 'NMFBAI', '美国:ISM:服务业PMI:就业': 'NMFEI', '美国:ISM:服务业PMI:新订单': 'NMFNOI', '美国:ISM:服务业PMI:物价': 'NMFPI', '美国:ISM:服务业PMI:商业活动': 'NMFBAI', '美国:ISM:服务业PMI:供应商交付': 'NMFSDI', '美国:ISM:服务业PMI:库存': 'NMFII', '美国:ISM:服务业PMI:订单库存': 'NMFOBI', '美国:ISM:服务业PMI:新出口订单': 'NMFEXI', '美国:ISM:服务业PMI:进口': 'NMFIMI', 'pmi_us_svc': 'NMFBAI', 'employment_us_季调': 'PAYEMS', # ── GDP ── 'gdp_us_yoy': 'A191RL1Q225SBEA', 'gdp_japan_yoy': 'JPNRGDPEXP', 'gdp_germany_yoy': 'CLVMNACSCAB1GQDE', 'gdp_uk_yoy': 'CLVMNACSCAB1GQUK', 'gdp_australia_yoy': 'NAEXKP01AUQ189S', # ── PPI/CPI ── 'ppi_us_yoy': 'PPIACO', 'ppi_us_同比': 'PPIACO', 'cpi_japan': 'JPNCPIALLMINMEI', 'ppi_japan': 'JPNPROINDMISMEI', 'ppi_eurozone': 'EA19PRMNIG01GYSAM', 'ipi_europe': 'EA19PRINTO01GYSAM', 'ppi_russia': 'RUSCPIALLMINMEI', # ── 欧洲/日本宏观 ── 'pmi_eu_mfg': 'BSCICP03EZM460S', 'ez_econ_sentiment': 'CSCICP03EZM460S', 'ez_consumer_conf': 'CSCICP03EZM460S', 'pmi_japan': 'BSCICP03JPM460S', 'pmi_canada': 'BSCICP03CAM460S', # ── OFR ── '美国:OFR金融压力指数': 'STLFSI2', # ── 美联储资产负债表 (FRED有汇总) ── 'fed_balance_sheet_total': 'WALCL', # ── 国际利率 (LIBOR→SOFR替代) ── 'libor_3m_英镑LIBOR3M': 'GBP3MTD156N', 'libor_3m_日元LIBOR3M': 'JPY3MTD156N', 'eurlibor_3m': 'EUR3MTD156N', # ── 美元指数(新兴市场) ── 'usd_index_em': 'DTWEXEMEGS', # ── 美国石油消费/需求 (FRED镜像EIA) ── 'us_oil_consumption': 'MTTIMUS2', 'demand_us': 'MTTIMUS2', # ── 美国贸易 ── 'us_crude_export_yoy': 'MCREXUS2', 'us_crude_import_yoy': 'MCRIMUS2', # ── 美联储资产负债表明细 (已验证可用) ── 'fed_balance_sheet_存款机构储备(不包括FHLB存款)': 'WRESBAL', 'fed_balance_sheet_证券回购协议(包括官方储备中的其他债权)': 'WLRRAL', 'fed_balance_sheet_对联邦政府的可支票存款': 'WDFOL', 'fed_balance_sheet_联邦储备银行股票': 'TREAST', 'fed_balance_sheet_联邦住房贷款银行(FHLB)借款': 'WSHOMCB', # ── 储备/需求 ── 'reserve_us': 'TOTRESNS', # ── HIBOR ── 'hibor_3m': 'HKONTD156N', # ── 中国PMI (FRED OECD proxy) ── 'pmi_china': 'CHNPMINDMISMEI', } def fetch_fred_bulk(start='2000-01-01'): """批量从 FRED 拉取所有映射的指标。""" from fredapi import Fred fred = Fred(api_key=FRED_KEY) end = datetime.now().strftime('%Y-%m-%d') results = {} total = len(FRED_MAP) print(f"\n[FRED] 拉取 {total} 个指标...") for i, (name, sid) in enumerate(FRED_MAP.items()): try: data = fred.get_series(sid, observation_start=start, observation_end=end) if data is not None and len(data) > 0: monthly = data.resample('ME').last().dropna() results[name] = monthly print(f" [{i+1}/{total}] ✓ {name:<35} {len(monthly)} 月") else: print(f" [{i+1}/{total}] ✗ {name:<35} 无数据") except Exception as e: print(f" [{i+1}/{total}] ✗ {name:<35} {str(e)[:40]}") df = pd.DataFrame(results) if results else pd.DataFrame() print(f" FRED 完成: {len(results)}/{total}") return df # ═══════════════════════════════════════════════════════════ # 2. YAHOO FINANCE — 商品/汇率/波动率 # ═══════════════════════════════════════════════════════════ YFINANCE_MAP = { # === 汇率 === 'usd_cny': 'CNY=X', 'usd_jpy': 'JPY=X', # === 波动率 === 'vix_nasdaq': '^VXN', 'ovx_crude_vol': '^OVX', '纳斯达克100波动率指数': '^VXN', # === 贵金属 === 'iron_ore_spot': 'TIOC.SI', 'nickel_price': '^SPGSNKP', '现货价(伦敦市场):黄金:美元': 'GC=F', '现货价(伦敦市场):白银:美元': 'SI=F', 'COMEX微型黄金': 'MGC=F', # === 大宗商品(替代 World Bank)=== '全球:实际市场价格:铜': 'HG=F', '全球:实际市场价格:铝': 'ALI=F', '全球:实际市场价格:锌': 'ZINC.L', '全球:实际市场价格:玉米': 'ZC=F', '全球:实际市场价格:大米': 'ZR=F', '全球:实际市场价格:棉花': 'CT=F', '全球:实际市场价格:豆粕': 'ZM=F', '全球:实际市场价格:豆油': 'ZL=F', '全球:实际市场价格:糖(CESE第11号合同最近的未来头寸)': 'SB=F', '全球:实际市场价格:天然气(印尼出口日本)': 'TTF=F', '全球:商品价格:原油:三大市场平均值': 'CL=F', # === 农产品 === 'agri_prices_大豆': 'ZS=F', 'agri_prices_小麦': 'ZW=F', 'agri_prices_燕麦': 'ZO=F', # === 更多商品(替代 World Bank 剩余)=== 'natgas_index': 'NG=F', # 天然气期货 '全球:实际市场价格:牛肉': 'LE=F', # 活牛期货 '全球:实际市场价格:橄榄油': 'ZL=F', # 豆油近似 '全球:实际市场价格:鱼粉': 'ZM=F', # 豆粕近似 # === 汇率proxy(用于供给/能源衍生)=== 'micex_brent': 'USDRUB=X', # USD/RUB作为MICEX-Brent proxy } def fetch_yfinance(start='2000-01-01'): """从 Yahoo Finance 拉取价格/汇率/商品数据(替代 World Bank)。""" import yfinance as yf tickers = {k: v for k, v in YFINANCE_MAP.items() if v is not None} total = len(tickers) print(f"\n[YFinance] 拉取 {total} 个指标(含商品替代 World Bank)...") results = {} for i, (name, ticker) in enumerate(tickers.items()): try: data = yf.download(ticker, start=start, progress=False, auto_adjust=True) if data is not None and len(data) > 0: close = data['Close'] if hasattr(close, 'columns'): close = close.iloc[:, 0] monthly = close.resample('ME').last().dropna() results[name] = monthly print(f" [{i+1}/{total}] ✓ {name:<35} {len(monthly)} 月") else: print(f" [{i+1}/{total}] ✗ {name:<35} 无数据") except Exception as e: print(f" [{i+1}/{total}] ✗ {name:<35} {str(e)[:40]}") df = pd.DataFrame(results) if results else pd.DataFrame() print(f" YFinance 完成: {len(results)}/{total}") return df # ═══════════════════════════════════════════════════════════ # 3. (已合并到 YFinance) World Bank 商品 → yfinance 期货替代 # ═══════════════════════════════════════════════════════════ def fetch_worldbank(): """World Bank 已被 yfinance 商品期货替代,返回空 DataFrame。""" print(f"\n[World Bank] 已合并到 YFinance 商品期货,跳过") return pd.DataFrame() # ═══════════════════════════════════════════════════════════ # 4. AKSHARE — 中国宏观 # ═══════════════════════════════════════════════════════════ def fetch_akshare(): """从 akshare 拉取中国宏观数据(M1/M2/PMI/GDP/CPI/PPI)。""" print(f"\n[akshare] 拉取中国宏观数据...") results = {} try: import akshare as ak # 中国 M1/M2 (货币供应量) try: ms = ak.macro_china_money_supply() if ms is not None and len(ms) > 0: ms['date'] = pd.to_datetime(ms.iloc[:, 0], errors='coerce') ms = ms.dropna(subset=['date']).set_index('date').sort_index() for col in ms.columns: ccl = col.lower() if 'M2' in col and '同比' in col: results['中国:M2:同比'] = pd.to_numeric(ms[col], errors='coerce') elif 'M1' in col and '同比' in col: results['中国:M1:同比'] = pd.to_numeric(ms[col], errors='coerce') elif 'M2' in col and '数量' in col: results['中国:M2'] = pd.to_numeric(ms[col], errors='coerce') elif 'M1' in col and '数量' in col: results['中国:M1'] = pd.to_numeric(ms[col], errors='coerce') for k in ['中国:M1', '中国:M2', '中国:M1:同比', '中国:M2:同比']: if k in results: print(f" ✓ {k:<25} {results[k].notna().sum()} 月") except Exception as e: print(f" ✗ money_supply: {str(e)[:40]}") # 中国 PMI try: pmi = ak.macro_china_pmi() if pmi is not None and len(pmi) > 0: pmi['date'] = pd.to_datetime(pmi.iloc[:, 0], errors='coerce') pmi = pmi.dropna(subset=['date']).set_index('date').sort_index() for col in pmi.columns: if '制造业' in col and '指数' in col: results['pmi_china_mfg'] = pd.to_numeric(pmi[col], errors='coerce') print(f" ✓ pmi_china_mfg {results['pmi_china_mfg'].notna().sum()} 月") elif '非制造业' in col and '指数' in col: results['pmi_china_nonsvc'] = pd.to_numeric(pmi[col], errors='coerce') print(f" ✓ pmi_china_nonsvc {results['pmi_china_nonsvc'].notna().sum()} 月") except Exception as e: print(f" ✗ pmi_china: {str(e)[:40]}") # 中国 CPI (月度) try: cpi = ak.macro_china_cpi_monthly() if cpi is not None and len(cpi) > 0: cpi['date'] = pd.to_datetime(cpi.iloc[:, 0], errors='coerce') cpi = cpi.dropna(subset=['date']).set_index('date').sort_index() if len(cpi.columns) > 1: results['cpi_china'] = pd.to_numeric(cpi.iloc[:, 0], errors='coerce') print(f" ✓ cpi_china {results['cpi_china'].notna().sum()} 月") except Exception as e: print(f" ✗ cpi_china: {str(e)[:40]}") # 中国 PPI try: ppi = ak.macro_china_ppi() if ppi is not None and len(ppi) > 0: ppi['date'] = pd.to_datetime(ppi.iloc[:, 0], errors='coerce') ppi = ppi.dropna(subset=['date']).set_index('date').sort_index() if len(ppi.columns) > 1: results['ppi_china_yoy'] = pd.to_numeric(ppi.iloc[:, 1], errors='coerce') print(f" ✓ ppi_china_yoy {results['ppi_china_yoy'].notna().sum()} 月") except Exception as e: print(f" ✗ ppi_china: {str(e)[:40]}") # 中国 GDP try: gdp = ak.macro_china_gdp() if gdp is not None and len(gdp) > 0: gdp['date'] = pd.to_datetime(gdp.iloc[:, 0], errors='coerce') gdp = gdp.dropna(subset=['date']).set_index('date').sort_index() for col in gdp.columns: if '同比增长' in col: results['gdp_china_yoy'] = pd.to_numeric(gdp[col], errors='coerce') results['gdp_china_growth_yoy'] = results['gdp_china_yoy'] print(f" ✓ gdp_china_yoy {results['gdp_china_yoy'].notna().sum()} 季") break except Exception as e: print(f" ✗ gdp_china: {str(e)[:40]}") # SHIBOR 3M try: shibor = ak.rate_interbank(market="上海银行同业拆借市场", symbol="Shibor人民币", indicator="3月") if shibor is not None and len(shibor) > 0: shibor['date'] = pd.to_datetime(shibor.iloc[:, 0], errors='coerce') shibor = shibor.set_index('date').sort_index() series = pd.to_numeric(shibor.iloc[:, 0], errors='coerce') monthly = series.resample('ME').last().dropna() results['shibor_3m'] = monthly print(f" ✓ shibor_3m {len(monthly)} 月") except Exception as e: print(f" ✗ shibor_3m: {str(e)[:40]}") except ImportError: print(" ✗ akshare 未安装 (pip install akshare)") except Exception as e: print(f" ✗ akshare 错误: {e}") df = pd.DataFrame(results) if results else pd.DataFrame() print(f" akshare 完成: {len(results)} 指标") return df # ═══════════════════════════════════════════════════════════ # 5. CFTC 持仓报告(自动下载) # ═══════════════════════════════════════════════════════════ def fetch_cftc(): """从 CFTC 官网自动下载持仓报告。""" print(f"\n[CFTC] 下载持仓报告...") results = {} try: # CFTC Disaggregated Futures-Only reports year = datetime.now().year url = f"https://www.cftc.gov/dea/newcot/{year}fut.zip" import zipfile, io resp = requests.get(url, timeout=30) if resp.status_code == 200: z = zipfile.ZipFile(io.BytesIO(resp.content)) fname = z.namelist()[0] df_cftc = pd.read_csv(z.open(fname)) # Filter for crude oil oil = df_cftc[df_cftc['Market_and_Exchange_Names'].str.contains( 'CRUDE OIL', case=False, na=False)] if len(oil) > 0: oil['date'] = pd.to_datetime(oil['As_of_Date_In_Form_YYMMDD'], format='%y%m%d', errors='coerce') oil = oil.set_index('date').sort_index() # Extract key positioning columns col_map = { 'ice_wti_mm_long': 'M_Money_Positions_Long_All', 'ice_wti_mm_short': 'M_Money_Positions_Short_All', } for name, col in col_map.items(): if col in oil.columns: series = pd.to_numeric(oil[col], errors='coerce') monthly = series.resample('ME').last().dropna() results[name] = monthly print(f" ✓ {name:<30} {len(monthly)} 月") print(f" CFTC 完成: {len(results)} 指标") else: print(f" ✗ CFTC 下载失败: HTTP {resp.status_code}") except Exception as e: print(f" ✗ CFTC 错误: {str(e)[:50]}") return pd.DataFrame(results) if results else pd.DataFrame() # ═══════════════════════════════════════════════════════════ # 6. EIA API # ═══════════════════════════════════════════════════════════ EIA_SERIES = { 'us_oil_inventory_total': 'PET.WCESTUS1.W', 'us_crude_production': 'PET.WCRFPUS2.W', } def fetch_eia(): """从 EIA API 拉取能源数据。""" print(f"\n[EIA] 拉取 {len(EIA_SERIES)} 个指标...") results = {} for name, sid in EIA_SERIES.items(): try: url = f"https://api.eia.gov/v2/seriesid/{sid}" resp = requests.get(url, params={'api_key': EIA_KEY}, timeout=30) if resp.status_code == 200: data = resp.json() if 'response' in data and 'data' in data['response']: records = data['response']['data'] df_tmp = pd.DataFrame(records) df_tmp['date'] = pd.to_datetime(df_tmp['period']) df_tmp['value'] = pd.to_numeric(df_tmp['value'], errors='coerce') series = df_tmp.set_index('date')['value'].sort_index() monthly = series.resample('ME').last().dropna() results[name] = monthly print(f" ✓ {name:<30} {len(monthly)} 月") else: print(f" ✗ {name:<30} 无数据") else: print(f" ✗ {name:<30} HTTP {resp.status_code}") except Exception as e: print(f" ✗ {name:<30} {str(e)[:40]}") return pd.DataFrame(results) if results else pd.DataFrame() # ═══════════════════════════════════════════════════════════ # 7. GPR 地缘政治风险指数 (Caldara & Iacoviello) # ═══════════════════════════════════════════════════════════ def fetch_gpr(): """从学术网站自动下载 GPR 地缘政治风险指数。 替代30个二值地缘事件虚拟变量。 来源: https://www.matteoiacoviello.com/gpr.htm """ print(f"\n[GPR] 下载地缘政治风险指数 (Caldara & Iacoviello)...") results = {} urls = [ "https://www.matteoiacoviello.com/gpr_files/data_gpr_monthly.xls", "https://www.matteoiacoviello.com/gpr_files/data_gpr_daily_recent.xls", ] try: df_gpr = None for url in urls: try: df_gpr = pd.read_excel(url) print(f" 下载成功: {url.split('/')[-1]}") break except Exception: continue if df_gpr is not None and len(df_gpr) > 0: # 找到日期列和GPR列 date_col = [c for c in df_gpr.columns if 'date' in c.lower() or 'month' in c.lower()] gpr_cols = [c for c in df_gpr.columns if 'gpr' in c.lower()] if date_col: df_gpr['date'] = pd.to_datetime(df_gpr[date_col[0]], errors='coerce') else: df_gpr['date'] = pd.to_datetime(df_gpr.iloc[:, 0], errors='coerce') df_gpr = df_gpr.dropna(subset=['date']).set_index('date').sort_index() for col in gpr_cols: series = pd.to_numeric(df_gpr[col], errors='coerce') monthly = series.resample('ME').last().dropna() if len(monthly) > 0: clean_name = col.strip().lower().replace(' ', '_') if clean_name in ('gpr', 'gprd'): clean_name = 'gpr_index' results[clean_name] = monthly print(f" ✓ {clean_name:<30} {len(monthly)} 月") if 'gpr_index' not in results and results: first_key = list(results.keys())[0] results['gpr_index'] = results[first_key] print(f" ✓ gpr_index (alias of {first_key})") else: print(f" ✗ GPR 数据为空") except Exception as e: print(f" ✗ GPR 下载失败: {str(e)[:50]}") df = pd.DataFrame(results) if results else pd.DataFrame() print(f" GPR 完成: {len(results)} 指标") return df # ═══════════════════════════════════════════════════════════ # MERGE + FEATURE ENGINEERING # ═══════════════════════════════════════════════════════════ def merge_all_sources(fred_df, yf_df, wb_df, ak_df, cftc_df, eia_df, gpr_df=None, existing_path='output/panel_monthly.csv'): """合并所有数据源,生成更新后的面板。""" print(f"\n[MERGE] 合并所有数据源...") existing = pd.read_csv(existing_path, index_col=0, parse_dates=True) print(f" 现有面板: {len(existing)} 行 × {len(existing.columns)} 列") print(f" 截止: {existing.index[-1].strftime('%Y-%m')}") sources = [ ('FRED', fred_df), ('YFinance', yf_df), ('WorldBank', wb_df), ('akshare', ak_df), ('CFTC', cftc_df), ('EIA', eia_df), ('GPR', gpr_df), ] updated_cols = 0 new_rows = 0 # 价格列需要用 FRED 数据覆盖已有值,确保数据源一致性 PRICE_OVERRIDE_COLS = {'WTI_spot', 'Brent_spot'} for src_name, src_df in sources: if src_df is None or len(src_df) == 0: continue for col in src_df.columns: if col in existing.columns: is_price_override = (col in PRICE_OVERRIDE_COLS and src_name == 'FRED') # 补全缺失值 + 追加新日期 for d in src_df.index: if d not in existing.index: existing.loc[d] = np.nan new_rows += 1 if pd.notna(src_df.at[d, col]): if pd.isna(existing.at[d, col]) or is_price_override: existing.at[d, col] = src_df.at[d, col] updated_cols += 1 else: # 新列 existing = existing.join(src_df[[col]], how='outer') existing = existing.sort_index() # 重新计算派生特征 print(f" 重新计算派生特征...") try: from data_pipeline import engineer_features existing = engineer_features(existing) except Exception as e: print(f" ⚠ engineer_features 跳过: {e}") # ISM 环比(22个):从 ISM 子指数自动计算 MoM ism_mom_pairs = [ ('美国:ISM:制造业PMI:环比', '美国:ISM:制造业PMI'), ('美国:ISM:制造业PMI:新订单:环比', '美国:ISM:制造业PMI:新订单'), ('美国:ISM:制造业PMI:订单库存:环比', '美国:ISM:制造业PMI:自有库存'), ('美国:ISM:制造业PMI:自有库存:环比', '美国:ISM:制造业PMI:自有库存'), ('美国:ISM:制造业PMI:客户库存:环比', '美国:ISM:制造业PMI:客户库存'), ('美国:ISM:制造业PMI:新出口订单:环比', '美国:ISM:制造业PMI:新出口订单'), ('美国:ISM:制造业PMI:进口:环比', '美国:ISM:制造业PMI:进口'), ('美国:ISM:制造业PMI:产出:环比', '美国:ISM:制造业PMI:产出'), ('美国:ISM:制造业PMI:供应商交付:环比', '美国:ISM:制造业PMI:供应商交付'), ('美国:ISM:制造业PMI:物价:环比', '美国:ISM:制造业PMI:物价'), ('美国:ISM:制造业PMI:就业:环比', '美国:ISM:制造业PMI:就业'), ('美国:ISM:服务业PMI:环比', '美国:ISM:服务业PMI'), ('美国:ISM:服务业PMI:商业活动:环比', '美国:ISM:服务业PMI:商业活动'), ('美国:ISM:服务业PMI:新订单:环比', '美国:ISM:服务业PMI:新订单'), ('美国:ISM:服务业PMI:就业:环比', '美国:ISM:服务业PMI:就业'), ('美国:ISM:服务业PMI:供应商交付:环比', '美国:ISM:服务业PMI:供应商交付'), ('美国:ISM:服务业PMI:库存:环比', '美国:ISM:服务业PMI:库存'), ('美国:ISM:服务业PMI:物价:环比', '美国:ISM:服务业PMI:物价'), ('美国:ISM:服务业PMI:订单库存:环比', '美国:ISM:服务业PMI:订单库存'), ('美国:ISM:服务业PMI:新出口订单:环比', '美国:ISM:服务业PMI:新出口订单'), ('美国:ISM:服务业PMI:进口:环比', '美国:ISM:服务业PMI:进口'), ('美国:ISM:服务业PMI:库存景气:环比', '美国:ISM:服务业PMI:库存'), ] ism_count = 0 for mom_col, base_col in ism_mom_pairs: if base_col in existing.columns: existing[mom_col] = existing[base_col].diff() ism_count += 1 if ism_count: print(f" ✓ ISM 环比: {ism_count} 个派生特征") out_path = os.path.join(OUTPUT_DIR, 'panel_monthly_live.csv') existing.to_csv(out_path) print(f" ✓ 更新后面板: {len(existing)} 行 × {len(existing.columns)} 列") print(f" ✓ 截止: {existing.index[-1].strftime('%Y-%m')}") print(f" ✓ 已保存: {out_path}") return existing # ═══════════════════════════════════════════════════════════ # TEST MODE # ═══════════════════════════════════════════════════════════ def test_apis(): """测试各 API 连通性。""" print("=" * 65) print("API 连通性测试") print("=" * 65) results = {} # FRED try: from fredapi import Fred fred = Fred(api_key=FRED_KEY) d = fred.get_series('DCOILWTICO', observation_start='2026-01-01') results['FRED'] = f"✓ 连通 ({len(d)} 条)" except Exception as e: results['FRED'] = f"✗ {str(e)[:40]}" # YFinance try: import yfinance as yf d = yf.download('CL=F', period='5d', progress=False) results['YFinance'] = f"✓ 连通 ({len(d)} 条)" except Exception as e: results['YFinance'] = f"✗ {str(e)[:40]}" # World Bank try: resp = requests.head("https://thedocs.worldbank.org/en/doc/5d903e848db1d1b83e0ec8f744e55570-0350012021/related/CMO-Historical-Data-Monthly.xlsx", timeout=10) results['WorldBank'] = f"✓ 可访问 (HTTP {resp.status_code})" except Exception as e: results['WorldBank'] = f"✗ {str(e)[:40]}" # akshare try: import akshare as ak results['akshare'] = f"✓ 已安装 (v{ak.__version__})" except Exception as e: results['akshare'] = f"✗ {str(e)[:40]}" # CFTC try: year = datetime.now().year resp = requests.head(f"https://www.cftc.gov/dea/newcot/{year}fut.zip", timeout=10) results['CFTC'] = f"✓ 可访问 (HTTP {resp.status_code})" except Exception as e: results['CFTC'] = f"✗ {str(e)[:40]}" # EIA try: resp = requests.get(f"https://api.eia.gov/v2/seriesid/PET.WCESTUS1.W", params={'api_key': EIA_KEY}, timeout=10) results['EIA'] = f"✓ 连通 (HTTP {resp.status_code})" except Exception as e: results['EIA'] = f"✗ {str(e)[:40]}" # GPR try: resp = requests.head("https://www.matteoiacoviello.com/gpr_files/data_gpr_monthly.xls", timeout=10) results['GPR'] = f"✓ 可访问 (HTTP {resp.status_code})" except Exception as e: results['GPR'] = f"✗ {str(e)[:40]}" print() for src, status in results.items(): print(f" {src:15s} {status}") ok = sum(1 for v in results.values() if '✓' in v) print(f"\n {ok}/{len(results)} 个数据源可用") return results # ═══════════════════════════════════════════════════════════ # MAIN # ═══════════════════════════════════════════════════════════ def main(test=False, source=None): print("=" * 65) print("全特征 API 数据更新") print("=" * 65) if test: test_apis() return start = '2000-01-01' # 按数据源拉取 fred_df = yf_df = wb_df = ak_df = cftc_df = eia_df = gpr_df = None if source is None or source == 'fred': fred_df = fetch_fred_bulk(start) if source is None or source == 'yfinance': yf_df = fetch_yfinance(start) if source is None or source == 'worldbank': wb_df = fetch_worldbank() if source is None or source == 'akshare': ak_df = fetch_akshare() if source is None or source == 'cftc': cftc_df = fetch_cftc() if source is None or source == 'eia': eia_df = fetch_eia() if source is None or source == 'gpr': gpr_df = fetch_gpr() # 合并 panel = merge_all_sources(fred_df, yf_df, wb_df, ak_df, cftc_df, eia_df, gpr_df) print(f"\n{'='*65}") print("✅ 数据更新完成") print(f"{'='*65}") if __name__ == '__main__': parser = argparse.ArgumentParser(description='全特征 API 数据更新') parser.add_argument('--test', action='store_true', help='测试 API 连通性') parser.add_argument('--source', choices=['fred', 'yfinance', 'worldbank', 'akshare', 'cftc', 'eia', 'gpr'], help='只更新某个数据源') args = parser.parse_args() main(test=args.test, source=args.source)