
Aj Stock Analysis
- 366 installs
- Updated July 27, 2026
- zuoa/aj-skills
Analyzes China A-share stocks with value-investing screening, deep single-stock analysis, sector comparison, and valuation using tushare public financial data.
About
A value-investing analysis tool for China A-share stocks that screens stocks, does deep single-stock and sector-comparison analysis, and computes valuation and intrinsic value from tushare financial data. A developer or investor uses it for low-frequency, fundamentals-based stock research.
- Value-investing screening, valuation, and financial-anomaly risk detection
- Uses tushare public financial data in a Python venv
Aj Stock Analysis by the numbers
- 366 all-time installs (skills.sh)
- Ranked #312 of 1,106 Finance & Trading skills by installs in the Skillselion catalog
- Data as of Jul 28, 2026 (Skillselion catalog sync)
npx skills add https://github.com/zuoa/aj-skills --skill aj-stock-analysisAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 366 |
|---|---|
| Last updated | July 27, 2026 |
| Repository | zuoa/aj-skills ↗ |
What it does
Analyzes China A-share stocks with value-investing screening, deep single-stock analysis, sector comparison, and valuation using tushare public financial data.
Files
China Stock Analysis Skill
基于价值投资理论的中国A股分析工具,面向低频交易的普通投资者。
When to Use
当用户请求以下操作时调用此skill:
- 分析某只A股股票
- 筛选符合条件的股票
- 对比多只股票或行业内股票
- 计算股票估值或内在价值
- 查看股票的财务健康状况
- 检测财务异常风险
Prerequisites
Python环境要求(必须使用venv)
所有脚本命令都应在项目虚拟环境中运行。
python3 -m venv .venv
source .venv/bin/activate安装依赖:
pip install tushare pandas numpyEnvironment Bootstrap(执行前必须自动完成)
在运行任何脚本前,先在当前项目目录执行:
cd <skill项目目录>
if [ ! -d ".venv" ]; then
python3 -m venv .venv
fi
source .venv/bin/activate
python3 -m pip install -U pip
python3 -m pip install tushare pandas numpy说明:
.venv必须位于 skill 项目根目录下(不是全局目录)- 若
import tushare失败,必须先执行上述 bootstrap,再继续后续分析流程
TUSHARE_TOKEN / BRAVE_API_KEY 建议保存在 ~/.aj-skills/.env,并在执行前先加载到当前 shell:
set -a
source ~/.aj-skills/.env
set +a说明:
- 脚本不会自动读取
~/.aj-skills/.env - 必须通过 CLI 参数显式传入:
--token "${TUSHARE_TOKEN}"、--brave-api-key "${BRAVE_API_KEY}"
执行前参数预检(使用 test 命令):
test -n "${TUSHARE_TOKEN}" || { echo "缺少 TUSHARE_TOKEN"; exit 1; }
# 仅当使用 Brave 新闻源时检查
test -n "${BRAVE_API_KEY}" || { echo "缺少 BRAVE_API_KEY(news-provider=brave 时必填)"; exit 1; }依赖检查
在执行任何分析前,先检查tushare是否已安装:
python3 -c "import tushare; print(tushare.__version__)"Core Modules
1. Stock Screener (股票筛选器)
筛选符合条件的股票
2. Financial Analyzer (财务分析器)
个股深度财务分析
3. Industry Comparator (行业对比)
同行业横向对比分析
4. Valuation Calculator (估值计算器)
内在价值测算与安全边际计算
5. News & Sentiment (新闻与舆情)
抓取近期社会面新闻并生成舆情风险评估
---
Workflow 1: Stock Screening (股票筛选)
用户请求筛选股票时使用。
Step 1: Collect Screening Criteria
向用户询问筛选条件。提供以下选项供用户选择或自定义:
估值指标:
- PE (市盈率): 例如 PE < 15
- PB (市净率): 例如 PB < 2
- PS (市销率): 例如 PS < 3
盈利能力:
- ROE (净资产收益率): 例如 ROE > 15%
- ROA (总资产收益率): 例如 ROA > 8%
- 毛利率: 例如 > 30%
- 净利率: 例如 > 10%
成长性:
- 营收增长率: 例如 > 10%
- 净利润增长率: 例如 > 15%
- 连续增长年数: 例如 >= 3年
股息:
- 股息率: 例如 > 3%
- 连续分红年数: 例如 >= 5年
财务安全:
- 资产负债率: 例如 < 60%
- 流动比率: 例如 > 1.5
- 速动比率: 例如 > 1
筛选范围:
- 全A股
- 沪深300成分股
- 中证500成分股
- 创业板/科创板
- 用户自定义列表
Step 2: Execute Screening
python scripts/stock_screener.py \
--scope "hs300" \
--token "${TUSHARE_TOKEN}" \
--pe-max 15 \
--roe-min 15 \
--debt-ratio-max 60 \
--dividend-min 2 \
--output screening_result.json参数说明:
--scope: 筛选范围 (all/hs300/zz500/cyb/kcb/custom:600519,000858,...)--pe-max/--pe-min: PE范围--pb-max/--pb-min: PB范围--roe-min: 最低ROE--growth-min: 最低增长率--debt-ratio-max: 最大资产负债率--dividend-min: 最低股息率--token: tushare token(必填)--format: 输出格式 (json/table)--quiet: 静默模式--output: 输出文件路径
Step 3: Present Results
读取 screening_result.json 并以表格形式呈现给用户:
| 代码 | 名称 | PE | PB | ROE | 股息率 | 评分 |
|---|---|---|---|---|---|---|
| 600519 | 贵州茅台 | 25.3 | 8.5 | 30.2% | 2.1% | 85 |
---
Workflow 2: Stock Analysis (个股分析)
用户请求分析某只股票时使用。
Step 1: Collect Stock Information
询问用户: 1. 股票代码或名称 2. 分析深度级别:
- 摘要级:关键指标 + 投资结论(1页)
- 标准级:财务分析 + 估值 + 行业对比 + 风险提示
- 深度级:完整调研报告,包含历史数据追踪
Step 1.5: Prepare Output Directory
单只股票分析时,skill 需要自动创建输出目录,命名规则:
${股票名称}_${股票代码}
示例:
stock_dir="贵州茅台_600519"
mkdir -p "${stock_dir}"Step 2: Fetch Stock Data
推荐使用“分模块抓取 + 聚合”流程(CLI解耦):
mkdir -p "${stock_dir}/data"
# 1) 基础信息
python scripts/fetch_basic.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--output "${stock_dir}/data/basic.json"
# 2) 财务数据
python scripts/fetch_financial.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--years 5 \
--output "${stock_dir}/data/financial.json"
# 3) 估值与行情
python scripts/fetch_valuation.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--output "${stock_dir}/data/valuation.json"
python scripts/fetch_price.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--days 180 \
--output "${stock_dir}/data/price.json"
# 4) 新闻舆情
python scripts/fetch_news_data.py \
--code "600519" \
--name "贵州茅台" \
--days 7 \
--limit 20 \
--provider brave \
--brave-api-key "${BRAVE_API_KEY}" \
--output "${stock_dir}/data/news.json"
# 5) 实时与事件窗口
python scripts/fetch_realtime.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--benchmark hs300 \
--window 60 \
--output "${stock_dir}/data/realtime.json"
python scripts/fetch_event_window.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--name "贵州茅台" \
--benchmark hs300 \
--pre-days 1 \
--post-days 1,3,5 \
--provider brave \
--brave-api-key "${BRAVE_API_KEY}" \
--output "${stock_dir}/data/event_window.json"
# 6) 聚合为分析输入
python scripts/assemble_data.py \
--input-dir "${stock_dir}/data" \
--output "${stock_dir}/stock_data.json"兼容模式(单命令抓取)保留如下:
python scripts/data_fetcher.py \
--code "600519" \
--token "${TUSHARE_TOKEN}" \
--data-type all \
--with-news \
--news-provider brave \
--brave-api-key "${BRAVE_API_KEY}" \
--with-realtime \
--with-event-window \
--benchmark hs300 \
--realtime-window 60 \
--event-window-pre 1 \
--event-window-post 1,3,5 \
--news-days 7 \
--news-limit 20 \
--years 5 \
--output "${stock_dir}/stock_data.json"参数说明:
--code: 股票代码--data-type: 数据类型 (basic/financial/valuation/holder/news/all)--years: 获取多少年的历史数据--token: tushare token(必填)--with-news: 附加新闻与舆情--news-days: 新闻窗口天数--news-limit: 新闻最大条数--news-sources: 新闻来源过滤(逗号分隔)--news-provider: 新闻源 (auto/brave/tushare/rss)--brave-api-key: Brave Search API Key(news-provider=brave时必填)--with-realtime: 附加实时指标(趋势/确认/风险/筹码)--with-event-window: 附加事件窗口分析(事件后1/3/5日反应)--benchmark: 相对强弱基准指数 (hs300/zz500/zz1000/cyb/kcb)--realtime-window: 实时指标计算窗口(日)--event-window-pre: 事件窗口前置天数--event-window-post: 事件窗口后验天数(逗号分隔)--cache-ttl-min: 缓存有效期(分钟)--format: 输出格式 (json/table)--quiet: 静默模式--output: 输出文件
可选:单独执行新闻舆情流程
python scripts/news_fetcher.py --code 600519 --name 贵州茅台 --token "${TUSHARE_TOKEN}" --days 7 --limit 20 --provider brave --brave-api-key "${BRAVE_API_KEY}" --output "${stock_dir}/news.json"
python scripts/sentiment_analyzer.py --input "${stock_dir}/news.json" --output "${stock_dir}/sentiment.json"Step 3: Run Financial Analysis
python scripts/financial_analyzer.py \
--input "${stock_dir}/stock_data.json" \
--level standard \
--output "${stock_dir}/analysis_result.json"或直接读取分模块目录(自动聚合):
python scripts/financial_analyzer.py \
--input-dir "${stock_dir}/data" \
--level standard \
--output "${stock_dir}/analysis_result.json"参数说明:
--input: 输入的股票数据文件--input-dir: 分模块目录(自动读取并聚合为分析输入)--level: 分析深度 (summary/standard/deep)--output: 输出文件
Step 4: Calculate Valuation
python scripts/valuation_calculator.py \
--input "${stock_dir}/stock_data.json" \
--methods dcf,ddm,relative \
--discount-rate 10 \
--growth-rate 8 \
--output "${stock_dir}/valuation_result.json"参数说明:
--input: 股票数据文件--methods: 估值方法 (dcf/ddm/relative/all)--discount-rate: 折现率(%)--terminal-growth: 永续增长率(%)--growth-rate: 永续增长率兼容别名(%)--margin-of-safety: 安全边际(%)--format: 输出格式 (json/table)--quiet: 静默模式--output: 输出文件
Step 5: Generate Report
读取分析结果,参考 templates/analysis_report.md 模板生成中文分析报告。
报告生成必检项(必须全部满足): 0. 最终报告必须落盘为 Markdown 文件(.md) 1. 必须包含“新闻与舆情”章节 2. 必须使用 stock_data.json 中的 news_sentiment/news_items 填充对应字段 3. 若新闻抓取失败,需在报告中明确写出失败原因(来自 news_sentiment.error) 4. 不允许省略模板中 summary_title 与“业绩与审计信号”章节 5. 若 stock_data.json 包含 realtime_metrics,报告必须包含“实时指标看板”章节(趋势/确认/风险/筹码) 6. 若 stock_data.json 包含 event_window,报告必须包含“事件窗口反应”内容(事件数、后1/3/5日收益、后1/3/5日超额收益) 7. 若 realtime_metrics 存在,综合评分必须采用 财务40% + 实时60%(实时优先)
报告结构(标准级): 1. 公司概况:基本信息、主营业务 2. 财务健康:资产负债表分析 3. 盈利能力:杜邦分析、利润率趋势 4. 成长性分析:营收/利润增长趋势 5. 实时指标看板:趋势/确认/风险/筹码 6. 事件窗口分析:事件后收益与超额收益 7. 估值分析:DCF/DDM/相对估值 8. 风险提示:财务异常检测、股东减持 9. 投资结论:综合评分、操作建议(实时优先)
报告标题规范:
summary_title使用格式:股票名称(股票代码):总结性结论- 示例:
贵州茅台(600519):财务稳健,估值与风险匹配度较好
输出文件:
${stock_dir}/final_report.mdStep 6: humanize Output
读取 ${stock_dir}/final_report.md,并调用 humanizer-zh skill 进行润色和优化。
输出文件:
${stock_dir}/final_report_humanized.md---
Workflow 3: Industry Comparison (行业对比)
CLI方式(板块分析,推荐)
# 1) 获取板块数据
python scripts/sector_fetcher.py \
--sector-name "算力板块" \
--token "${TUSHARE_TOKEN}" \
--sector-file config/sector_computing_default.json \
--output "${stock_dir}/sector_data.json"
# 2) 生成板块分析结果 + Markdown报告
python scripts/sector_analyze.py \
--input "${stock_dir}/sector_data.json" \
--output "${stock_dir}/sector_analysis.json"Step 1: Collect Comparison Targets
询问用户: 1. 目标股票代码(可多个) 2. 或者:行业分类 + 对比数量
Step 2: Fetch Industry Data
python scripts/data_fetcher.py \
--codes "600519,000858,002304" \
--token "${TUSHARE_TOKEN}" \
--data-type comparison \
--output industry_data.json或按行业获取:
python scripts/data_fetcher.py \
--industry "白酒" \
--token "${TUSHARE_TOKEN}" \
--top 10 \
--output industry_data.jsonStep 3: Generate Comparison
python scripts/financial_analyzer.py \
--input industry_data.json \
--mode comparison \
--output comparison_result.jsonStep 4: Present Comparison Table
| 指标 | 贵州茅台 | 五粮液 | 洋河股份 | 行业均值 |
|---|---|---|---|---|
| PE | 25.3 | 18.2 | 15.6 | 22.4 |
| ROE | 30.2% | 22.5% | 20.1% | 18.5% |
| 毛利率 | 91.5% | 75.2% | 72.3% | 65.4% |
| 评分 | 85 | 78 | 75 | - |
---
Workflow 4: Valuation Calculator (估值计算)
Step 1: Collect Valuation Parameters
询问用户估值参数(或使用默认值):
DCF模型参数:
- 折现率 (WACC): 默认10%
- 预测期: 默认5年
- 永续增长率: 默认3%
DDM模型参数:
- 要求回报率: 默认10%
- 股息增长率: 使用历史数据推算
相对估值参数:
- 对比基准: 行业均值 / 历史均值
Step 2: Run Valuation
python scripts/valuation_calculator.py \
--code "600519" \
--methods all \
--discount-rate 10 \
--terminal-growth 3 \
--forecast-years 5 \
--margin-of-safety 30 \
--output valuation.jsonStep 3: Present Valuation Results
| 估值方法 | 内在价值 | 当前价格 | 安全边际价格 | 结论 |
|---|---|---|---|---|
| DCF | ¥2,150 | ¥1,680 | ¥1,505 | 低估 |
| DDM | ¥1,980 | ¥1,680 | ¥1,386 | 低估 |
| 相对估值 | ¥1,850 | ¥1,680 | ¥1,295 | 合理 |
---
Financial Anomaly Detection (财务异常检测)
在分析过程中自动检测以下异常信号:
检测项目
1. 应收账款异常
- 应收账款增速 > 营收增速 × 1.5
- 应收账款周转天数大幅增加
2. 现金流背离
- 净利润持续增长但经营现金流下降
- 现金收入比 < 80%
3. 存货异常
- 存货增速 > 营收增速 × 2
- 存货周转天数大幅增加
4. 毛利率异常
- 毛利率波动 > 行业均值波动 × 2
- 毛利率与同行严重偏离
5. 关联交易
- 关联交易占比过高(> 30%)
6. 股东减持
- 大股东近期减持公告
- 高管集中减持
风险等级
- 🟢 低风险:无明显异常
- 🟡 中风险:1-2项轻微异常
- 🔴 高风险:多项异常或严重异常
---
A-Share Specific Analysis (A股特色分析)
政策敏感度
根据行业分类提供政策相关提示:
- 房地产:房住不炒政策
- 新能源:补贴政策变化
- 医药:集采政策影响
- 互联网:反垄断、数据安全
股东结构分析
1. 控股股东类型(国企/民企/外资) 2. 股权集中度 3. 近期增减持情况 4. 质押比例
---
Output Format
JSON/Table输出格式
- 默认
json - 可选
--format table用于终端快速查看 - 使用
--quiet可关闭过程日志
所有脚本输出JSON格式,便于后续处理:
{
"code": "600519",
"name": "贵州茅台",
"analysis_date": "2025-01-25",
"level": "standard",
"summary": {
"score": 85,
"conclusion": "低估",
"recommendation": "建议关注"
},
"financials": { ... },
"valuation": { ... },
"risks": [ ... ]
}Markdown报告
生成结构化的中文Markdown报告,参考 templates/analysis_report.md。
---
Data Contract
核心数据结构由 scripts/data_contract.py 约束。分析脚本会在运行前校验:
- 顶层必需字段:
code/fetch_time/data_type/basic_info - 常用可选字段:
financial_data/financial_indicators/valuation/price/holder/dividend - 新闻相关字段:
news_items/news_sentiment - 业绩审计字段:
performance_data(含forecast/express/audit/main_business) - 报表字段要求:
financial_data.balance_sheet必须是数组financial_data.income_statement必须是数组financial_data.cash_flow必须是数组
字段映射(Akshare -> Tushare)
| 兼容语义 | 当前字段(推荐) | 兼容别名/来源 |
|---|---|---|
| PE(TTM) | valuation.latest.pe_ttm | valuation.latest.pe |
| PB | valuation.latest.pb | - |
| 净利润 | financial_data.income_statement[].净利润 | n_income |
| 经营现金流净额 | financial_data.cash_flow[].经营活动产生的现金流量净额 | n_cashflow_act |
| 资本开支现金 | financial_data.cash_flow[].购建固定资产、无形资产和其他长期资产支付的现金 | c_pay_acq_const_fiolta |
| ROE | financial_indicators[].净资产收益率 | roe |
| 资产负债率 | financial_indicators[].资产负债率 | debt_to_assets |
---
Error Handling
网络错误
如果tushare数据获取失败,提示用户: 1. 检查网络连接 2. 稍后重试(可能是接口限流) 3. 尝试更换数据源
股票代码无效
提示用户检查股票代码是否正确,提供可能的匹配建议。
数据不完整
对于新上市股票或财务数据不完整的情况,说明数据限制并基于可用数据进行分析。
---
Best Practices
1. 数据时效性:财务数据以最新季报/年报为准,价格数据为当日收盘价;开启 --with-realtime / --with-event-window 时补充趋势与事件冲击的动态指标 2. 投资建议:所有分析仅供参考,不构成投资建议 3. 风险提示:始终包含风险提示,特别是财务异常检测结果 4. 对比分析:单只股票分析时,自动包含行业均值对比 5. 评分权重:实时数据存在时使用 财务40% + 实时60%;实时缺失时退化为财务分
Important Notes
- 所有分析基于公开财务数据,不涉及任何内幕信息
- 估值模型的参数假设对结果影响较大,需向用户说明
- A股市场受政策影响较大,定量分析需结合定性判断
# Tushare credential
TUSHARE_TOKEN=your_tushare_token_here
# Optional runtime defaults
STOCK_ANALYSIS_DEFAULT_YEARS=3
STOCK_ANALYSIS_DEFAULT_NEWS_DAYS=7
STOCK_ANALYSIS_DEFAULT_NEWS_LIMIT=20
{
"AI芯片": {
"688256": "寒武纪-U",
"688041": "海光信息",
"688047": "龙芯中科",
"300474": "景嘉微"
},
"服务器": {
"603019": "中科曙光",
"000977": "浪潮信息",
"601138": "工业富联",
"000938": "紫光股份"
},
"光模块/CPO": {
"300308": "中际旭创",
"300502": "新易盛",
"300394": "天孚通信",
"000988": "华工科技"
},
"IDC/数据中心": {
"603881": "数据港",
"300738": "奥飞数据",
"300383": "光环新网"
}
}
A股市场特色分析指南
一、A股市场特点
1.1 市场结构
交易所分布:
| 交易所 | 定位 | 特点 |
|---|---|---|
| 上交所主板 | 大型成熟企业 | 蓝筹为主 |
| 深交所主板 | 大中型企业 | 蓝筹为主 |
| 创业板 | 创新创业企业 | 高成长高风险 |
| 科创板 | 科技创新企业 | 注册制、门槛高 |
| 北交所 | 创新型中小企业 | 专精特新 |
投资者结构:
- 散户占比高(约60%交易量)
- 机构化程度提升中
- 公募、私募、外资(北向资金)
1.2 与成熟市场差异
| 特征 | A股 | 美股 |
|---|---|---|
| 涨跌幅限制 | 主板±10%,创/科±20% | 无限制 |
| T+1交易 | 是 | T+0 |
| 退市机制 | 较宽松 | 较严格 |
| 做空机制 | 有限 | 完善 |
| 分红文化 | 逐步加强 | 成熟 |
---
二、政策与行业分析
2.1 政策敏感行业
强监管行业:
| 行业 | 主要政策 | 影响 |
|---|---|---|
| 房地产 | "房住不炒"、三道红线 | 融资、销售双限 |
| 教育 | "双减"政策 | 学科类培训受限 |
| 互联网 | 反垄断、数据安全 | 业务合规成本 |
| 医药 | 集中采购 | 仿制药利润压缩 |
| 游戏 | 版号审批、防沉迷 | 发展受限 |
政策支持行业:
| 行业 | 政策方向 | 机会 |
|---|---|---|
| 新能源 | 碳中和、补贴 | 产业链高增长 |
| 半导体 | 国产替代 | 政策资金支持 |
| 高端制造 | 专精特新 | 税收优惠 |
| 医疗创新 | 创新药支持 | 审批加速 |
| 消费升级 | 扩大内需 | 长期机会 |
2.2 行业周期判断
周期性行业特征:
- 钢铁、煤炭、有色、化工、航运
- 盈利随经济周期大幅波动
- 低谷时PE可能极高,高峰时PE反而低
周期位置判断: 1. 产能利用率 2. 库存水平 3. 产品价格 4. 行业资本开支
周期股投资策略:
- 高PE时买入(周期底部)
- 低PE时卖出(周期顶部)
- 与普通估值逻辑相反
---
三、股东结构分析
3.1 控股股东类型
国有企业:
- 优势:资源获取、政策支持、稳定性
- 劣势:效率可能较低、激励不足
- 关注:混改进展、管理层变动
民营企业:
- 优势:机制灵活、激励充分
- 劣势:治理风险、资金链风险
- 关注:实控人信用、股权质押
外资控股:
- 优势:治理规范、技术先进
- 劣势:可能转移利润、政策风险
- 关注:关联交易、转让定价
3.2 股权结构关注点
股权集中度:
前十大股东持股比例
- > 70%:高度集中
- 50-70%:中度集中
- < 50%:较为分散股权质押:
质押比例 = 已质押股份 / 总持股
- < 20%:正常
- 20-50%:需关注
- > 50%:高风险减持公告:
- 计划减持:关注减持比例和原因
- 实际减持:跟踪减持进度
- 大宗交易:可能是减持信号
3.3 管理层分析
管理层持股:
- 有持股:利益一致
- 股权激励计划:行权条件是否合理
高管变动:
- 核心高管离职需警惕
- 董秘、财务总监变动尤需关注
- 审计师更换是红旗
---
四、财务风险识别
4.1 常见财务问题
收入确认激进:
信号:
- 应收账款增速 > 营收增速 × 1.5
- 第四季度收入占比过高
- 关联方销售占比高费用资本化:
信号:
- 开发支出/研发费用比例异常高
- 在建工程长期不转固
- 无形资产/商誉大幅增加存货问题:
信号:
- 存货周转率持续下降
- 存货跌价准备计提不充分
- 存货增速远超营收增速4.2 现金流验证
盈利质量检验:
# 现金收入比
现金收入比 = 销售商品收到的现金 / 营业收入
# 正常应 > 100%(含增值税)
# 净利润现金含量
净现比 = 经营活动现金流净额 / 净利润
# 正常应 > 80%异常模式:
| 净利润 | 经营现金流 | 含义 |
|---|---|---|
| 增长 | 增长 | 健康 |
| 增长 | 下降 | 警惕 |
| 下降 | 增长 | 可能改善 |
| 下降 | 下降 | 恶化 |
4.3 A股常见雷区
商誉减值:
- 并购形成大额商誉
- 被收购公司业绩不达预期
- 年末集中计提减值
应收账款坏账:
- 大客户集中
- 账龄结构恶化
- 关联方应收款
存货跌价:
- 产品价格下跌
- 库存积压
- 技术更新导致贬值
---
五、A股特有数据源
5.1 重要公告类型
| 公告类型 | 关注点 |
|---|---|
| 业绩预告/快报 | 业绩变动幅度 |
| 股东减持公告 | 减持比例和原因 |
| 股权激励计划 | 行权条件和覆盖范围 |
| 关联交易公告 | 交易定价公允性 |
| 对外担保公告 | 担保风险敞口 |
| 资产重组公告 | 重组方案和估值 |
| 诉讼仲裁公告 | 涉诉金额和进展 |
5.2 监管信息
交易所问询函:
- 年报问询函:关注重点问题
- 关注函:可能存在违规
- 监管函:存在明确问题
证监会公告:
- 立案调查
- 行政处罚
- 市场禁入
5.3 北向资金
北向资金数据:
- 每日净流入/流出
- 个股持仓变动
- 持股市值占比参考价值:
- 外资偏好稳定的白马股
- 大幅增减持可能有信号意义
- 但不宜盲目跟随
---
六、估值调整因素
6.1 A股估值特点
历史估值中枢:
- 沪深300 PE中枢:约12-15倍
- 创业板 PE中枢:约35-50倍
- 受流动性和情绪影响大
估值溢价/折价:
| 因素 | 溢价 | 折价 |
|---|---|---|
| 稀缺性 | 行业龙头、独特资源 | 同质化竞争 |
| 确定性 | 消费、医药 | 周期、概念 |
| 成长性 | 高增长预期 | 增长放缓 |
| 治理 | 良好治理历史 | 历史有问题 |
| 流动性 | 大市值、高换手 | 小市值、低流动性 |
6.2 A股估值方法调整
DCF调整:
- 折现率可能需要更高(政策不确定性)
- 永续增长率需谨慎(产业政策变化)
相对估值调整:
- 参考行业历史估值区间
- 考虑市场情绪周期位置
- 参考港股/美股可比公司
---
七、投资决策框架
7.1 A股价值投资检查清单
基本面筛选:
- [ ] ROE连续5年 > 15%
- [ ] 毛利率稳定或上升
- [ ] 经营现金流/净利润 > 0.8
- [ ] 资产负债率 < 60%(非金融)
- [ ] 营收和利润正增长
风险排查:
- [ ] 无重大诉讼/违规
- [ ] 无异常关联交易
- [ ] 控股股东质押率 < 50%
- [ ] 无频繁减持公告
- [ ] 审计意见为标准无保留
估值判断:
- [ ] PE < 行业均值
- [ ] PB处于历史低位
- [ ] 存在明确的安全边际
政策评估:
- [ ] 行业政策无重大利空
- [ ] 公司受益于政策方向
- [ ] 合规风险可控
7.2 买入/卖出时机
买入时机考量: 1. 估值处于历史低位(如PE分位数 < 30%) 2. 市场悲观情绪过度 3. 基本面稳健,短期利空不改长期逻辑 4. 安全边际 > 30%
卖出时机考量: 1. 估值处于历史高位(如PE分位数 > 70%) 2. 基本面出现实质性恶化 3. 政策环境发生重大不利变化 4. 发现更好的投资机会
7.3 持仓管理
A股仓位建议:
- 单只股票:不超过20%
- 单一行业:不超过30%
- 保持一定现金比例应对波动
定期检视:
- 每季度财报后重新评估
- 关注重大公告
- 跟踪行业政策变化
财务指标详解
一、估值指标
1.1 市盈率 (PE - Price to Earnings)
计算公式:
PE = 股价 / 每股收益 (EPS)
PE = 总市值 / 净利润类型:
- 静态PE:使用上一年度净利润
- 动态PE:使用预期未来净利润
- TTM PE:使用过去12个月滚动净利润
解读标准:
| PE范围 | 通常含义 |
|---|---|
| < 10 | 可能低估或有问题 |
| 10-15 | 相对便宜 |
| 15-25 | 合理估值 |
| 25-40 | 高成长预期 |
| > 40 | 可能高估 |
注意事项:
- 周期股低谷时PE可能极高
- 亏损公司PE无意义
- 需要与行业均值比较
1.2 市净率 (PB - Price to Book)
计算公式:
PB = 股价 / 每股净资产
PB = 总市值 / 净资产(股东权益)解读标准:
| PB范围 | 通常含义 |
|---|---|
| < 1 | 可能低估(破净) |
| 1-2 | 相对合理 |
| 2-5 | 高于账面价值 |
| > 5 | 高度依赖无形资产 |
适用场景:
- 重资产行业:银行、地产、钢铁
- 周期底部估值
- 清算价值参考
1.3 市销率 (PS - Price to Sales)
计算公式:
PS = 股价 / 每股销售额
PS = 总市值 / 营业收入适用场景:
- 亏损但高增长公司
- 互联网、SaaS企业
- 同行业对比
1.4 企业价值倍数 (EV/EBITDA)
计算公式:
EV = 总市值 + 总负债 - 现金
EBITDA = 息税折旧摊销前利润
EV/EBITDA = 企业价值 / EBITDA优势:
- 剔除资本结构差异
- 适合跨国比较
- 适用于高杠杆行业
---
二、盈利能力指标
2.1 净资产收益率 (ROE - Return on Equity)
计算公式:
ROE = 净利润 / 平均净资产 × 100%杜邦分析分解:
ROE = 净利率 × 资产周转率 × 权益乘数
= (净利润/营收) × (营收/总资产) × (总资产/净资产)解读标准:
| ROE范围 | 评价 |
|---|---|
| < 8% | 较低,需关注原因 |
| 8-15% | 一般水平 |
| 15-25% | 优秀 |
| > 25% | 卓越(需验证可持续性) |
Buffett标准: 连续5年ROE > 15%
2.2 总资产收益率 (ROA - Return on Assets)
计算公式:
ROA = 净利润 / 平均总资产 × 100%与ROE对比:
- ROA 反映资产运营效率
- 高杠杆可能导致 ROE高但ROA一般
- 银行等行业ROA通常较低
2.3 毛利率 (Gross Profit Margin)
计算公式:
毛利率 = (营收 - 营业成本) / 营收 × 100%行业参考:
| 行业 | 典型毛利率 |
|---|---|
| 白酒 | 70-90% |
| 医药 | 50-80% |
| 软件 | 60-80% |
| 零售 | 20-40% |
| 制造业 | 15-30% |
2.4 净利率 (Net Profit Margin)
计算公式:
净利率 = 净利润 / 营收 × 100%关注点:
- 毛利率高但净利率低 → 期间费用高
- 净利率波动大 → 非经常性损益多
---
三、财务安全指标
3.1 资产负债率 (Debt Ratio)
计算公式:
资产负债率 = 总负债 / 总资产 × 100%标准参考:
| 资产负债率 | 评价 |
|---|---|
| < 40% | 保守 |
| 40-60% | 适中 |
| 60-70% | 偏高 |
| > 70% | 高风险(非金融行业) |
注意:
- 银行业通常 > 90%,属正常
- 需区分有息负债和经营性负债
3.2 流动比率 (Current Ratio)
计算公式:
流动比率 = 流动资产 / 流动负债标准:
- 一般要求 > 1.5
- < 1 表示短期偿债压力大
3.3 速动比率 (Quick Ratio)
计算公式:
速动比率 = (流动资产 - 存货) / 流动负债标准:
- 一般要求 > 1
- 剔除存货变现不确定性
3.4 利息覆盖倍数 (Interest Coverage)
计算公式:
利息覆盖倍数 = EBIT / 利息费用标准:
- 一般要求 > 3
- < 1.5 存在偿债风险
---
四、运营效率指标
4.1 应收账款周转率/天数
计算公式:
应收周转率 = 营收 / 平均应收账款
应收周转天数 = 365 / 应收周转率解读:
- 周转天数越短越好
- 行业差异大,需对比同行
- 持续延长需警惕
4.2 存货周转率/天数
计算公式:
存货周转率 = 营业成本 / 平均存货
存货周转天数 = 365 / 存货周转率行业参考:
- 零售业:30-60天
- 制造业:60-120天
- 白酒:可能超过1年(正常)
4.3 总资产周转率
计算公式:
总资产周转率 = 营收 / 平均总资产解读:
- 反映资产利用效率
- 轻资产公司通常较高
- 杜邦分析重要组成部分
---
五、成长性指标
5.1 营收增长率
计算公式:
营收增长率 = (本期营收 - 上期营收) / 上期营收 × 100%
复合增长率 CAGR = (末期/初期)^(1/年数) - 15.2 净利润增长率
计算公式:
净利润增长率 = (本期净利 - 上期净利) / 上期净利 × 100%注意:
- 区分扣非净利润与归母净利润
- 排除非经常性损益影响
5.3 每股收益增长率 (EPS Growth)
重要性:
- 直接关系到股东回报
- 需考虑股本变动影响
---
六、股东回报指标
6.1 股息率 (Dividend Yield)
计算公式:
股息率 = 每股股息 / 股价 × 100%A股参考:
- 银行股:4-6%
- 公用事业:3-5%
- 消费品:2-4%
- 平均水平:约2%
6.2 分红比例 (Payout Ratio)
计算公式:
分红比例 = 现金分红总额 / 净利润 × 100%解读:
- 30-50% 较为健康
- 过高可能影响发展
- 过低可能不重视股东回报
6.3 股息增长率
计算公式:
股息CAGR = (末期股息/初期股息)^(1/年数) - 1用于DDM估值模型
---
七、现金流指标
7.1 经营现金流
重要性:
- 反映真实盈利质量
- 比净利润更难操纵
关注点:
现金收入比 = 销售收到的现金 / 营收- 通常应 > 100%
7.2 自由现金流 (FCF)
计算公式:
FCF = 经营现金流净额 - 资本支出解读:
- 正FCF:可用于分红、偿债、回购
- 持续负FCF:需要外部融资
7.3 现金流与净利润对比
净利润/经营现金流 比值:
- < 0.8:盈利质量可疑
- 0.8-1.2:正常
- > 1.2:现金回收好于账面
---
八、财务异常信号
8.1 红旗警示
| 信号 | 可能含义 |
|---|---|
| 应收账款增速 >> 营收增速 | 收入确认激进 |
| 存货增速 >> 营收增速 | 产品滞销或囤积 |
| 净利润增长但现金流下降 | 盈利质量差 |
| 毛利率异常高于同行 | 可能虚增收入 |
| 关联交易占比过高 | 利益输送风险 |
| 频繁更换审计师 | 财务问题风险 |
| 经营现金流持续为负 | 商业模式存疑 |
8.2 Beneish M-Score 简化版
核心检测变量: 1. DSRI:应收账款天数变化 2. GMI:毛利率变化 3. AQI:资产质量变化 4. SGI:销售增长 5. TATA:应计项目/总资产
M-Score > -1.78 时,财务造假概率较高
价值投资核心原则
一、价值投资理论基础
1.1 Benjamin Graham 的核心理念
安全边际 (Margin of Safety)
- 以显著低于内在价值的价格买入
- 建议安全边际至少 25-30%
- 安全边际越大,投资风险越低
内在价值 (Intrinsic Value)
- 公司真实价值由未来现金流决定
- 独立于市场价格波动
- 需要通过财务分析计算
Mr. Market 寓言
- 市场是情绪化的"市场先生"
- 价格常常偏离价值
- 利用市场波动而非被其左右
1.2 Warren Buffett 的投资哲学
护城河 (Economic Moat)
- 品牌优势:如茅台、可口可乐
- 成本优势:规模经济、独特资源
- 网络效应:用户越多价值越大
- 转换成本:客户迁移成本高
- 政府牌照:特许经营权
能力圈原则
- 只投资自己能理解的业务
- 简单的商业模式更容易分析
- 复杂不等于好
长期持有
- "我们最喜欢的持有期是永远"
- 复利效应需要时间发挥
- 减少交易成本和税负
---
二、Graham 经典选股标准
2.1 防御型投资者标准
| 指标 | 标准 | 说明 |
|---|---|---|
| 公司规模 | 年销售额 > 1亿美元 | 避免小公司风险 |
| 财务稳健 | 流动比率 > 2 | 短期偿债能力 |
| 盈利稳定 | 过去10年每年盈利 | 商业模式验证 |
| 分红历史 | 连续20年派息 | 股东回报意识 |
| 盈利增长 | 10年EPS增长 > 33% | 年均 3% 以上 |
| PE比率 | PE < 15 | 避免估值泡沫 |
| PB比率 | PB < 1.5 | 资产保护 |
| PE × PB | < 22.5 | 综合估值限制 |
2.2 积极型投资者标准
在防御型基础上,可接受:
- 规模较小但增长更快的公司
- PE略高但ROE持续优秀
- 行业地位突出的公司
---
三、估值方法论
3.1 DCF 现金流折现模型
基本公式:
内在价值 = Σ (FCFt / (1+r)^t) + 终值自由现金流 (FCF) 计算:
FCF = 经营现金流 - 资本支出关键参数:
- 折现率 (r):通常 8-12%
- 增长率:分阶段估计
- 终值:永续增长模型或退出倍数
适用场景:
- 现金流稳定可预测的公司
- 成熟期企业
- 公用事业、消费品等
3.2 DDM 股息折现模型
Gordon 增长模型:
内在价值 = D1 / (r - g)其中:
- D1 = 下一年预期股息
- r = 要求回报率
- g = 股息永续增长率
适用场景:
- 高分红企业
- 银行、保险、公用事业
- 分红政策稳定的公司
3.3 相对估值法
市盈率 (PE) 法:
合理价格 = EPS × 合理PE参考标准:
- 行业平均PE
- 历史平均PE
- 可比公司PE
市净率 (PB) 法:
合理价格 = 每股净资产 × 合理PB适用于:
- 重资产行业
- 银行等金融机构
- 周期性行业低谷期
---
四、选股检查清单
4.1 定量筛选标准
估值合理:
- [ ] PE < 行业均值
- [ ] PB < 3(非轻资产行业)
- [ ] PE × PB < 22.5
盈利能力强:
- [ ] ROE > 15%(连续3年)
- [ ] 毛利率 > 行业均值
- [ ] 净利率稳定或上升
财务安全:
- [ ] 资产负债率 < 60%
- [ ] 流动比率 > 1.5
- [ ] 利息覆盖倍数 > 3
成长性:
- [ ] 营收年复合增长 > 10%
- [ ] 净利润年复合增长 > 10%
4.2 定性评估要点
商业模式:
- [ ] 业务是否容易理解
- [ ] 是否有护城河
- [ ] 行业竞争格局
管理层:
- [ ] 管理层持股情况
- [ ] 过往诚信记录
- [ ] 资本配置能力
行业前景:
- [ ] 行业增长空间
- [ ] 政策支持/限制
- [ ] 技术变革风险
---
五、风险控制原则
5.1 分散投资
- 单只股票仓位不超过 20%
- 行业分散:不同行业至少 3-5 个
- 避免高度相关的股票
5.2 仓位管理
根据安全边际调整仓位:
- 安全边际 > 50%:可重仓
- 安全边际 30-50%:正常仓位
- 安全边际 < 30%:轻仓或观望
5.3 卖出原则
应该卖出的情况: 1. 公司基本面恶化 2. 股价严重高估(超出内在价值 50%+) 3. 发现更好的投资机会 4. 当初买入的理由不再成立
不应该卖出的理由: 1. 仅仅因为股价下跌 2. 短期市场波动 3. 宏观经济担忧(除非影响公司基本面)
---
六、A股价值投资特点
6.1 与成熟市场的差异
| 特点 | A股 | 美股 |
|---|---|---|
| 投资者结构 | 散户为主 | 机构为主 |
| 市场效率 | 相对较低 | 相对较高 |
| 估值波动 | 较大 | 相对平稳 |
| 政策影响 | 显著 | 较小 |
| 价值投资空间 | 可能更大 | 更成熟 |
6.2 A股价值投资机会
低效市场机会:
- 市场情绪极端时往往出现低估
- 中小盘股可能被机构忽视
- 行业轮动带来的错杀机会
需要额外关注:
- 财务报表真实性
- 大股东行为
- 政策变化影响
- 再融资和减持
#!/usr/bin/env python3
"""
分模块数据聚合器:
- 读取各模块 JSON 文件
- 合并为 financial_analyzer 可用的 stock_data.json
"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
from typing import Dict, Optional
SECTION_KEYS = [
"basic_info",
"financial_data",
"financial_indicators",
"performance_data",
"valuation",
"price",
"flow_metrics",
"chip_events",
"realtime_metrics",
"event_window",
"holder",
"dividend",
"news_items",
"news_sentiment",
]
def _load_json(path: Path) -> Dict:
return json.loads(path.read_text(encoding="utf-8"))
def _extract_code(payload: Dict) -> str:
code = payload.get("code")
if code:
return str(code)
basic = payload.get("basic_info")
if isinstance(basic, dict) and basic.get("code"):
return str(basic.get("code"))
return ""
def _merge_one(target: Dict, payload: Dict, section_name: str, strict_code: bool = True) -> None:
src_code = _extract_code(payload)
if src_code:
if target.get("code") and target.get("code") != src_code and strict_code:
raise ValueError(f"代码不一致: target={target.get('code')} section({section_name})={src_code}")
target["code"] = target.get("code") or src_code
for key in SECTION_KEYS:
if key in payload and payload[key] is not None:
target[key] = payload[key]
sections = target.setdefault("_sections", {})
sections[section_name] = {
"fetch_time": payload.get("fetch_time"),
"source_file": section_name,
}
def assemble_stock_data(
code: str = "",
basic_file: Optional[Path] = None,
financial_file: Optional[Path] = None,
valuation_file: Optional[Path] = None,
price_file: Optional[Path] = None,
holder_file: Optional[Path] = None,
news_file: Optional[Path] = None,
realtime_file: Optional[Path] = None,
event_window_file: Optional[Path] = None,
strict_code: bool = True,
) -> Dict:
result: Dict = {
"code": str(code or ""),
"fetch_time": datetime.now().isoformat(),
"data_type": "assembled",
"basic_info": {},
}
file_map = {
"basic": basic_file,
"financial": financial_file,
"valuation": valuation_file,
"price": price_file,
"holder": holder_file,
"news": news_file,
"realtime": realtime_file,
"event_window": event_window_file,
}
for section_name, section_path in file_map.items():
if section_path is None:
continue
path = Path(section_path)
if not path.exists():
continue
payload = _load_json(path)
if not isinstance(payload, dict):
raise ValueError(f"文件内容不是 object: {path}")
_merge_one(result, payload, section_name=section_name, strict_code=strict_code)
result.setdefault("_section_files", {})[section_name] = str(path)
if not result.get("code"):
raise ValueError("聚合后缺少 code,请至少提供一个包含 code 的 section 文件或 --code")
if not isinstance(result.get("basic_info"), dict):
result["basic_info"] = {}
if not result["basic_info"].get("code"):
result["basic_info"]["code"] = result["code"]
return result
def assemble_from_dir(input_dir: str, code: str = "", strict_code: bool = True) -> Dict:
root = Path(input_dir)
return assemble_stock_data(
code=code,
basic_file=(root / "basic.json"),
financial_file=(root / "financial.json"),
valuation_file=(root / "valuation.json"),
price_file=(root / "price.json"),
holder_file=(root / "holder.json"),
news_file=(root / "news.json"),
realtime_file=(root / "realtime.json"),
event_window_file=(root / "event_window.json"),
strict_code=strict_code,
)
def _resolve_default_file(input_dir: Optional[str], filename: str) -> Optional[Path]:
if not input_dir:
return None
path = Path(input_dir) / filename
return path if path.exists() else None
def main():
parser = argparse.ArgumentParser(description="聚合分模块股票数据为 stock_data.json")
parser.add_argument("--code", default="", help="股票代码(可选,缺失时从 section 文件推断)")
parser.add_argument("--input-dir", default="", help="模块文件目录(默认文件名: basic.json 等)")
parser.add_argument("--basic-file", default="", help="basic 文件路径")
parser.add_argument("--financial-file", default="", help="financial 文件路径")
parser.add_argument("--valuation-file", default="", help="valuation 文件路径")
parser.add_argument("--price-file", default="", help="price 文件路径")
parser.add_argument("--holder-file", default="", help="holder 文件路径")
parser.add_argument("--news-file", default="", help="news 文件路径")
parser.add_argument("--realtime-file", default="", help="realtime 文件路径")
parser.add_argument("--event-window-file", default="", help="event_window 文件路径")
parser.add_argument("--no-strict-code", action="store_true", help="关闭代码一致性校验")
parser.add_argument("--output", required=True, help="输出 stock_data.json 路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
basic_file = Path(args.basic_file) if args.basic_file else _resolve_default_file(args.input_dir, "basic.json")
financial_file = Path(args.financial_file) if args.financial_file else _resolve_default_file(args.input_dir, "financial.json")
valuation_file = Path(args.valuation_file) if args.valuation_file else _resolve_default_file(args.input_dir, "valuation.json")
price_file = Path(args.price_file) if args.price_file else _resolve_default_file(args.input_dir, "price.json")
holder_file = Path(args.holder_file) if args.holder_file else _resolve_default_file(args.input_dir, "holder.json")
news_file = Path(args.news_file) if args.news_file else _resolve_default_file(args.input_dir, "news.json")
realtime_file = Path(args.realtime_file) if args.realtime_file else _resolve_default_file(args.input_dir, "realtime.json")
event_window_file = (
Path(args.event_window_file)
if args.event_window_file
else _resolve_default_file(args.input_dir, "event_window.json")
)
if args.input_dir and not any([args.basic_file, args.financial_file, args.valuation_file, args.price_file, args.holder_file, args.news_file, args.realtime_file, args.event_window_file]):
result = assemble_from_dir(
input_dir=args.input_dir,
code=args.code,
strict_code=not args.no_strict_code,
)
else:
result = assemble_stock_data(
code=args.code,
basic_file=basic_file,
financial_file=financial_file,
valuation_file=valuation_file,
price_file=price_file,
holder_file=holder_file,
news_file=news_file,
realtime_file=realtime_file,
event_window_file=event_window_file,
strict_code=not args.no_strict_code,
)
output_path = Path(args.output)
output_path.parent.mkdir(parents=True, exist_ok=True)
output_path.write_text(json.dumps(result, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"聚合完成: {output_path}")
print(f"code={result.get('code')} sections={list((result.get('_section_files') or {}).keys())}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
统一数据契约定义与校验工具。
"""
from typing import Dict, List, Optional, Tuple
REQUIRED_TOP_LEVEL_KEYS = [
"code",
"fetch_time",
"data_type",
"basic_info",
]
OPTIONAL_SECTIONS = [
"financial_data",
"performance_data",
"financial_indicators",
"valuation",
"price",
"flow_metrics",
"chip_events",
"realtime_metrics",
"event_window",
"holder",
"dividend",
"news_items",
"news_sentiment",
]
def validate_stock_data(data: Dict, required_sections: Optional[List[str]] = None) -> Tuple[bool, List[str]]:
"""校验股票数据结构是否满足分析脚本要求。"""
errors: List[str] = []
if not isinstance(data, dict):
return False, ["输入数据不是 JSON object"]
for key in REQUIRED_TOP_LEVEL_KEYS:
if key not in data:
errors.append(f"缺少顶层字段: {key}")
if "basic_info" in data and not isinstance(data.get("basic_info"), dict):
errors.append("字段 basic_info 必须是 object")
if "financial_data" in data:
financial_data = data.get("financial_data")
if not isinstance(financial_data, dict):
errors.append("字段 financial_data 必须是 object")
else:
for key in ["balance_sheet", "income_statement", "cash_flow"]:
if key in financial_data and not isinstance(financial_data.get(key), list):
errors.append(f"字段 financial_data.{key} 必须是 array")
for key in [
"financial_indicators",
"news_items",
"holder",
"dividend",
"valuation",
"price",
"news_sentiment",
"performance_data",
"flow_metrics",
"chip_events",
"realtime_metrics",
"event_window",
]:
if key in data and data[key] is not None:
expected_type = list if key in ["financial_indicators", "news_items"] else dict
if not isinstance(data[key], expected_type):
type_name = "array" if expected_type is list else "object"
errors.append(f"字段 {key} 必须是 {type_name}")
if "performance_data" in data and isinstance(data.get("performance_data"), dict):
perf = data["performance_data"]
for key in ["forecast", "express", "audit"]:
if key in perf and not isinstance(perf.get(key), list):
errors.append(f"字段 performance_data.{key} 必须是 array")
if required_sections:
for section in required_sections:
if section not in data:
errors.append(f"缺少分析所需字段: {section}")
return len(errors) == 0, errors
def ensure_stock_data(data: Dict, required_sections: Optional[List[str]] = None) -> None:
"""校验失败时抛出 ValueError。"""
ok, errors = validate_stock_data(data, required_sections=required_sections)
if not ok:
raise ValueError("数据结构校验失败: " + "; ".join(errors))
#!/usr/bin/env python3
"""
A股数据获取模块
使用tushare获取股票财务数据、行情数据、股东信息等
依赖: pip install tushare pandas
"""
import argparse
import json
import os
import sys
import time
from datetime import datetime, timedelta
from functools import wraps
from typing import Callable, Dict, Optional, Sequence, Tuple
from news_fetcher import fetch_news
from sentiment_analyzer import analyze_news_sentiment
from realtime_metrics import calculate_realtime_metrics
from event_window import collect_event_candidates, calculate_event_window
try:
import pandas as pd
import tushare as ts
except ImportError:
print("错误: 请先安装依赖库")
print("pip install tushare pandas")
sys.exit(1)
INDEX_CODE_MAP = {
"hs300": "000300.SH",
"zz500": "000905.SH",
"zz1000": "000852.SH",
"cyb": "399006.SZ",
"kcb": "000688.SH",
}
def retry_on_failure(max_retries: int = 3, delay: float = 1.0):
"""网络请求重试装饰器"""
def decorator(func: Callable):
@wraps(func)
def wrapper(*args, **kwargs):
last_error = None
for attempt in range(max_retries):
try:
return func(*args, **kwargs)
except Exception as exc:
last_error = exc
if attempt < max_retries - 1:
time.sleep(delay * (attempt + 1))
return {"error": f"重试{max_retries}次后失败: {str(last_error)}"}
return wrapper
return decorator
def get_tushare_pro(token: str):
"""初始化tushare pro客户端"""
if not token:
raise RuntimeError("缺少 tushare token,请通过 --token 传入")
ts.set_token(token)
return ts.pro_api()
PRO = None
CLI_TOKEN = None
VERBOSE = True
def pro():
"""延迟初始化tushare客户端。"""
global PRO
if PRO is None:
PRO = get_tushare_pro(CLI_TOKEN)
return PRO
def log(message: str):
if VERBOSE:
print(message)
def format_table(headers, rows) -> str:
widths = [len(str(h)) for h in headers]
for row in rows:
for i, cell in enumerate(row):
widths[i] = max(widths[i], len(str(cell)))
sep = "-+-".join("-" * w for w in widths)
lines = [" | ".join(str(h).ljust(widths[i]) for i, h in enumerate(headers)), sep]
for row in rows:
lines.append(" | ".join(str(cell).ljust(widths[i]) for i, cell in enumerate(row)))
return "\n".join(lines)
def safe_float(value) -> Optional[float]:
"""安全转换为浮点数"""
if value is None or value == "" or value == "--":
return None
try:
if pd.isna(value):
return None
if isinstance(value, str):
value = value.replace("%", "").replace(",", "")
return float(value)
except (ValueError, TypeError):
return None
def parse_event_window_post_days(raw: str) -> Tuple[int, ...]:
"""解析事件窗口后验天数参数,如: 1,3,5"""
if not raw:
return (1, 3, 5)
values = []
for part in str(raw).split(","):
part = part.strip()
if not part:
continue
try:
day = int(part)
if day >= 1:
values.append(day)
except Exception:
continue
if not values:
return (1, 3, 5)
return tuple(sorted(set(values)))
def get_cache_path(code: str, data_type: str) -> str:
"""获取缓存文件路径(可按TTL过期)"""
cache_dir = os.path.join(os.path.dirname(__file__), ".cache")
os.makedirs(cache_dir, exist_ok=True)
return os.path.join(cache_dir, f"{code}_{data_type}.json")
def get_legacy_cache_path(code: str, data_type: str) -> str:
"""兼容旧版按日缓存文件命名"""
cache_dir = os.path.join(os.path.dirname(__file__), ".cache")
os.makedirs(cache_dir, exist_ok=True)
today = datetime.now().strftime("%Y%m%d")
return os.path.join(cache_dir, f"{code}_{data_type}_{today}.json")
def _parse_iso_dt(value: str) -> Optional[datetime]:
if not value:
return None
try:
return datetime.fromisoformat(value.replace("Z", "+00:00")).replace(tzinfo=None)
except Exception:
return None
def load_cache(code: str, data_type: str, ttl_minutes: int = 1440) -> Optional[dict]:
"""加载缓存数据(按TTL有效)"""
paths = [get_cache_path(code, data_type), get_legacy_cache_path(code, data_type)]
for cache_path in paths:
if not os.path.exists(cache_path):
continue
try:
with open(cache_path, "r", encoding="utf-8") as f:
data = json.load(f)
if ttl_minutes and ttl_minutes > 0:
ts = _parse_iso_dt(data.get("_cache_saved_at", "")) or _parse_iso_dt(data.get("fetch_time", ""))
if ts:
age_min = (datetime.now() - ts).total_seconds() / 60
if age_min > ttl_minutes:
continue
return data
except (json.JSONDecodeError, IOError):
continue
return None
def save_cache(code: str, data_type: str, data: dict):
"""保存缓存数据"""
cache_path = get_cache_path(code, data_type)
try:
to_save = dict(data)
to_save["_cache_saved_at"] = datetime.now().isoformat()
with open(cache_path, "w", encoding="utf-8") as f:
json.dump(to_save, f, ensure_ascii=False, default=str)
except IOError:
pass
def normalize_symbol(code: str) -> str:
"""统一为6位symbol"""
code = (code or "").strip().upper()
if "." in code:
return code.split(".")[0]
return code
def to_ts_code(code: str) -> str:
"""将symbol转换为ts_code"""
code = (code or "").strip().upper()
if "." in code:
return code
symbol = normalize_symbol(code)
if symbol.startswith(("6", "9", "5")):
suffix = "SH"
elif symbol.startswith(("4", "8")):
suffix = "BJ"
else:
suffix = "SZ"
return f"{symbol}.{suffix}"
def ts_to_symbol(ts_code: str) -> str:
"""ts_code转symbol"""
return (ts_code or "").split(".")[0]
def latest_trade_date() -> str:
"""获取最近交易日"""
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=30)).strftime("%Y%m%d")
cal = pro().trade_cal(exchange="", start_date=start_date, end_date=end_date, is_open="1")
if cal is None or cal.empty:
return end_date
cal = cal.sort_values("cal_date")
return str(cal.iloc[-1]["cal_date"])
def _format_date_ymd(date_str: str) -> str:
"""YYYYMMDD -> YYYY-MM-DD"""
if not date_str or len(str(date_str)) != 8:
return date_str
s = str(date_str)
return f"{s[:4]}-{s[4:6]}-{s[6:8]}"
@retry_on_failure(max_retries=2, delay=1.0)
def get_stock_info(code: str) -> dict:
"""获取股票基本信息"""
try:
ts_code = to_ts_code(code)
basic = pro().stock_basic(
ts_code=ts_code,
list_status="L",
fields="ts_code,symbol,name,industry,list_date",
)
if basic is None or basic.empty:
return {"code": normalize_symbol(code), "error": "未找到股票基础信息"}
info = basic.iloc[0].to_dict()
db = pro().daily_basic(
ts_code=ts_code,
fields="ts_code,trade_date,total_mv,circ_mv,total_share,float_share,pe_ttm,pb",
limit=1,
)
latest_db = db.iloc[0].to_dict() if db is not None and not db.empty else {}
return {
"code": info.get("symbol", normalize_symbol(code)),
"name": info.get("name", ""),
"industry": info.get("industry", ""),
"market_cap": safe_float(latest_db.get("total_mv")) * 10000 if latest_db.get("total_mv") is not None else None,
"float_cap": safe_float(latest_db.get("circ_mv")) * 10000 if latest_db.get("circ_mv") is not None else None,
"total_shares": safe_float(latest_db.get("total_share")) * 10000 if latest_db.get("total_share") is not None else None,
"float_shares": safe_float(latest_db.get("float_share")) * 10000 if latest_db.get("float_share") is not None else None,
"pe_ttm": safe_float(latest_db.get("pe_ttm")),
"pb": safe_float(latest_db.get("pb")),
"listing_date": _format_date_ymd(info.get("list_date", "")),
"ts_code": ts_code,
}
except Exception as exc:
return {"code": normalize_symbol(code), "error": str(exc)}
def _df_to_records(df: pd.DataFrame, max_records: int) -> list:
if df is None or df.empty:
return []
if "end_date" in df.columns:
df = df.sort_values("end_date", ascending=False)
return df.head(max_records).to_dict(orient="records")
def _add_income_aliases(records: list) -> list:
for row in records:
row["净利润"] = row.get("n_income")
row["营业收入"] = row.get("revenue")
row["主营业务收入"] = row.get("total_revenue")
return records
def _add_cashflow_aliases(records: list) -> list:
for row in records:
row["经营活动产生的现金流量净额"] = row.get("n_cashflow_act")
row["购建固定资产、无形资产和其他长期资产支付的现金"] = row.get("c_pay_acq_const_fiolta")
return records
def _add_balance_aliases(records: list) -> list:
for row in records:
row["总资产"] = row.get("total_assets")
row["总负债"] = row.get("total_liab")
assets = safe_float(row.get("total_assets"))
liab = safe_float(row.get("total_liab"))
row["资产负债率"] = (liab / assets * 100) if assets and liab is not None else None
return records
@retry_on_failure(max_retries=2, delay=1.0)
def get_financial_data(code: str, years: int = 3) -> dict:
"""获取财务数据(资产负债表、利润表、现金流量表)"""
max_records = min(years * 4, 12)
ts_code = to_ts_code(code)
result = {
"balance_sheet": [],
"income_statement": [],
"cash_flow": [],
}
try:
bs = pro().balancesheet(ts_code=ts_code, limit=max_records)
result["balance_sheet"] = _add_balance_aliases(_df_to_records(bs, max_records))
except Exception as exc:
result["balance_sheet_error"] = str(exc)
try:
inc = pro().income(ts_code=ts_code, limit=max_records)
result["income_statement"] = _add_income_aliases(_df_to_records(inc, max_records))
except Exception as exc:
result["income_statement_error"] = str(exc)
try:
cf = pro().cashflow(ts_code=ts_code, limit=max_records)
result["cash_flow"] = _add_cashflow_aliases(_df_to_records(cf, max_records))
except Exception as exc:
result["cash_flow_error"] = str(exc)
return result
def _map_indicator_row(row: Dict) -> Dict:
"""tushare指标映射为现有中文字段"""
return {
**row,
"日期": row.get("end_date"),
"净资产收益率": row.get("roe"),
"加权净资产收益率": row.get("roe_waa"),
"总资产报酬率": row.get("roa"),
"销售毛利率": row.get("grossprofit_margin"),
"销售净利率": row.get("netprofit_margin"),
"资产负债率": row.get("debt_to_assets"),
"流动比率": row.get("current_ratio"),
"速动比率": row.get("quick_ratio"),
"应收账款周转率": row.get("arturn"),
"应收账款周转天数": row.get("ar_days"),
"存货周转率": row.get("invturn"),
"存货周转天数": row.get("inv_days"),
"总资产周转率": row.get("assets_turn"),
"营业收入增长率": row.get("tr_yoy"),
"主营业务收入增长率": row.get("tr_yoy"),
"净利润增长率": row.get("netprofit_yoy"),
"应收账款增长率": row.get("recp_yoy"),
"存货增长率": row.get("inv_yoy"),
"权益乘数": row.get("assets_to_eqt"),
}
def get_financial_indicators(code: str, limit: int = 8) -> list:
"""获取财务指标"""
ts_code = to_ts_code(code)
try:
df = pro().fina_indicator(ts_code=ts_code, limit=limit)
if df is None or df.empty:
return []
if "end_date" in df.columns:
df = df.sort_values("end_date", ascending=False)
records = df.head(limit).to_dict(orient="records")
return [_map_indicator_row(r) for r in records]
except Exception:
return []
@retry_on_failure(max_retries=2, delay=1.0)
def get_performance_data(code: str, years: int = 3) -> dict:
"""获取业绩预告/快报/审计意见/主营构成。"""
ts_code = to_ts_code(code)
start_date = (datetime.now() - timedelta(days=365 * max(1, years))).strftime("%Y%m%d")
end_date = datetime.now().strftime("%Y%m%d")
result = {
"forecast": [],
"express": [],
"audit": [],
"main_business": {
"by_product": [],
"by_region": [],
},
}
try:
df = pro().forecast(ts_code=ts_code, start_date=start_date, end_date=end_date)
if df is not None and not df.empty:
df = df.sort_values("ann_date", ascending=False)
result["forecast"] = df.head(12).to_dict(orient="records")
except Exception as exc:
result["forecast_error"] = str(exc)
try:
df = pro().express(ts_code=ts_code, start_date=start_date, end_date=end_date)
if df is not None and not df.empty:
df = df.sort_values("ann_date", ascending=False)
result["express"] = df.head(8).to_dict(orient="records")
except Exception as exc:
result["express_error"] = str(exc)
try:
df = pro().fina_audit(ts_code=ts_code, start_date=start_date, end_date=end_date)
if df is not None and not df.empty:
df = df.sort_values("ann_date", ascending=False)
result["audit"] = df.head(8).to_dict(orient="records")
except Exception as exc:
result["audit_error"] = str(exc)
for biz_type, key in [("P", "by_product"), ("D", "by_region")]:
try:
df = pro().fina_mainbz(ts_code=ts_code, type=biz_type, start_date=start_date, end_date=end_date)
if df is not None and not df.empty:
# 仅保留最近一期构成并按收入排序
latest_period = str(df["end_date"].max())
latest_df = df[df["end_date"] == latest_period].copy()
if "bz_sales" in latest_df.columns:
latest_df["bz_sales"] = pd.to_numeric(latest_df["bz_sales"], errors="coerce")
latest_df = latest_df.sort_values("bz_sales", ascending=False)
result["main_business"][key] = latest_df.head(20).to_dict(orient="records")
except Exception as exc:
result[f"main_business_{key}_error"] = str(exc)
return result
def get_valuation_data(code: str) -> dict:
"""获取估值数据"""
ts_code = to_ts_code(code)
result = {}
try:
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=3 * 365)).strftime("%Y%m%d")
df = pro().daily_basic(
ts_code=ts_code,
start_date=start_date,
end_date=end_date,
fields="ts_code,trade_date,pe_ttm,pb",
)
if df is None or df.empty:
return result
df = df.sort_values("trade_date")
latest = df.iloc[-1].to_dict()
result["latest"] = {
"date": latest.get("trade_date"),
"pe_ttm": safe_float(latest.get("pe_ttm")),
"pb": safe_float(latest.get("pb")),
}
result["history_count"] = len(df)
pe = safe_float(latest.get("pe_ttm"))
pb = safe_float(latest.get("pb"))
if pe is not None:
series = pd.to_numeric(df["pe_ttm"], errors="coerce").dropna()
if not series.empty:
result["pe_percentile"] = float((series <= pe).mean() * 100)
if pb is not None:
series = pd.to_numeric(df["pb"], errors="coerce").dropna()
if not series.empty:
result["pb_percentile"] = float((series <= pb).mean() * 100)
except Exception as exc:
result["error"] = str(exc)
result["note"] = "估值历史数据获取失败,将使用基本信息中的估值"
return result
@retry_on_failure(max_retries=2, delay=1.0)
def get_holder_data(code: str) -> dict:
"""获取股东信息"""
ts_code = to_ts_code(code)
result = {}
try:
top10 = pro().top10_holders(ts_code=ts_code)
if top10 is not None and not top10.empty:
top10 = top10.sort_values(["end_date", "hold_ratio"], ascending=[False, False])
result["top_10_holders"] = top10.head(10).to_dict(orient="records")
except Exception as exc:
result["top_10_holders_error"] = str(exc)
try:
holders = pro().stk_holdernumber(ts_code=ts_code)
if holders is not None and not holders.empty:
holders = holders.sort_values("end_date", ascending=False)
result["holder_count_history"] = holders.head(10).to_dict(orient="records")
except Exception as exc:
result["holder_count_error"] = str(exc)
return result
@retry_on_failure(max_retries=2, delay=1.0)
def get_dividend_data(code: str) -> dict:
"""获取分红数据"""
ts_code = to_ts_code(code)
try:
df = pro().dividend(ts_code=ts_code)
if df is None or df.empty:
return {"dividend_history": [], "dividend_count": 0}
if "end_date" in df.columns:
df = df.sort_values("end_date", ascending=False)
records = df.to_dict(orient="records")
for row in records:
row["每股股利"] = row.get("cash_div_tax")
row["派息"] = row.get("cash_div_tax")
return {
"dividend_history": records,
"dividend_count": len(records),
}
except Exception:
return {"dividend_history": [], "dividend_count": 0}
@retry_on_failure(max_retries=2, delay=1.0)
def get_price_data(code: str, days: int = 60) -> dict:
"""获取价格数据"""
ts_code = to_ts_code(code)
try:
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=max(40, days * 3))).strftime("%Y%m%d")
df = pro().daily(
ts_code=ts_code,
start_date=start_date,
end_date=end_date,
fields="ts_code,trade_date,open,high,low,close,pct_chg,vol,amount",
)
if df is None or df.empty:
return {}
df = df.sort_values("trade_date").tail(max(60, days))
latest = df.iloc[-1]
db_map = {}
try:
db = pro().daily_basic(
ts_code=ts_code,
start_date=start_date,
end_date=end_date,
fields="ts_code,trade_date,turnover_rate,volume_ratio",
)
if db is not None and not db.empty:
db = db.sort_values("trade_date")
db_map = {
str(row.get("trade_date")): {
"turnover_rate": safe_float(row.get("turnover_rate")),
"volume_ratio": safe_float(row.get("volume_ratio")),
}
for _, row in db.iterrows()
}
except Exception:
db_map = {}
price_data = []
for _, row in df.tail(min(max(days, 60), 180)).iterrows():
trade_date = str(row.get("trade_date"))
ext = db_map.get(trade_date, {})
price_data.append(
{
"日期": trade_date,
"开盘": row.get("open"),
"最高": row.get("high"),
"最低": row.get("low"),
"收盘": row.get("close"),
"涨跌幅": row.get("pct_chg"),
"成交量": row.get("vol"),
"成交额": row.get("amount"),
"量比": ext.get("volume_ratio"),
"换手率": ext.get("turnover_rate"),
}
)
tr_20 = [safe_float(x.get("换手率")) for x in price_data[-20:]]
tr_20 = [x for x in tr_20 if x is not None]
return {
"latest_price": safe_float(latest.get("close")),
"latest_date": str(latest.get("trade_date")),
"price_change_pct": safe_float(latest.get("pct_chg")),
"volume": safe_float(latest.get("vol")),
"turnover": safe_float(latest.get("amount")),
"high_60d": safe_float(df["high"].max()),
"low_60d": safe_float(df["low"].min()),
"avg_volume_20d": safe_float(df.tail(20)["vol"].mean()),
"avg_amount_20d": safe_float(df.tail(20)["amount"].mean()),
"avg_turnover_rate_20d": safe_float(sum(tr_20) / len(tr_20)) if tr_20 else None,
"price_data": price_data,
}
except Exception as exc:
return {"error": str(exc)}
@retry_on_failure(max_retries=2, delay=1.0)
def get_index_price_data(index_name: str, days: int = 60) -> dict:
"""获取基准指数价格数据(用于相对强弱)"""
index_code = INDEX_CODE_MAP.get(index_name)
if not index_code:
return {"error": f"不支持的benchmark: {index_name}"}
try:
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=max(40, days * 3))).strftime("%Y%m%d")
df = pro().index_daily(
ts_code=index_code,
start_date=start_date,
end_date=end_date,
fields="ts_code,trade_date,open,high,low,close,pct_chg,vol,amount",
)
if df is None or df.empty:
return {}
df = df.sort_values("trade_date").tail(max(60, days))
latest = df.iloc[-1]
rows = []
for _, row in df.tail(min(max(days, 60), 180)).iterrows():
rows.append(
{
"日期": str(row.get("trade_date")),
"开盘": row.get("open"),
"最高": row.get("high"),
"最低": row.get("low"),
"收盘": row.get("close"),
"涨跌幅": row.get("pct_chg"),
"成交量": row.get("vol"),
"成交额": row.get("amount"),
}
)
return {
"benchmark": index_name,
"benchmark_ts_code": index_code,
"latest_price": safe_float(latest.get("close")),
"latest_date": str(latest.get("trade_date")),
"price_change_pct": safe_float(latest.get("pct_chg")),
"price_data": rows,
}
except Exception as exc:
return {"error": str(exc)}
def get_flow_metrics(code: str, days: int = 20) -> dict:
"""获取资金流指标(大单净流入及连续性)。"""
ts_code = to_ts_code(code)
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=max(20, days * 3))).strftime("%Y%m%d")
try:
df = pro().moneyflow(ts_code=ts_code, start_date=start_date, end_date=end_date)
if df is None or df.empty:
return {}
df = df.sort_values("trade_date").tail(max(5, days))
net_values = []
for _, row in df.iterrows():
buy_lg = safe_float(row.get("buy_lg_amount")) or 0.0
buy_elg = safe_float(row.get("buy_elg_amount")) or 0.0
sell_lg = safe_float(row.get("sell_lg_amount")) or 0.0
sell_elg = safe_float(row.get("sell_elg_amount")) or 0.0
net = buy_lg + buy_elg - sell_lg - sell_elg
net_values.append(net)
latest_row = df.iloc[-1]
latest_net = net_values[-1] if net_values else 0.0
latest_amount = safe_float(latest_row.get("amount"))
net_ratio = (latest_net / latest_amount * 100.0) if latest_amount not in [None, 0] else None
tail_5 = net_values[-5:] if len(net_values) >= 5 else net_values
positive_days_5 = len([x for x in tail_5 if x > 0])
return {
"latest_trade_date": str(latest_row.get("trade_date")),
"latest_net_inflow": round(latest_net, 4),
"net_inflow_ratio": round(net_ratio, 4) if net_ratio is not None else None,
"positive_days_5": positive_days_5,
}
except Exception as exc:
return {"error": str(exc)}
def get_chip_events(code: str, basic_info: Optional[dict] = None) -> dict:
"""获取筹码供给侧事件(解禁/减持/回购)。"""
ts_code = to_ts_code(code)
basic_info = basic_info or {}
float_shares = safe_float(basic_info.get("float_shares"))
total_shares = safe_float(basic_info.get("total_shares"))
end_date = datetime.now().strftime("%Y%m%d")
start_30 = (datetime.now() + timedelta(days=30)).strftime("%Y%m%d")
start_90 = (datetime.now() + timedelta(days=90)).strftime("%Y%m%d")
back_30 = (datetime.now() - timedelta(days=30)).strftime("%Y%m%d")
back_90 = (datetime.now() - timedelta(days=90)).strftime("%Y%m%d")
result = {
"unlock_30d_ratio": 0.0,
"unlock_90d_ratio": 0.0,
"reduction_density_30d": 0.0,
"repurchase_ratio_90d": 0.0,
}
try:
df = pro().share_float(ts_code=ts_code)
if df is not None and not df.empty:
df = df.sort_values("float_date")
df30 = df[(df["float_date"] <= start_30) & (df["float_date"] >= end_date)]
df90 = df[(df["float_date"] <= start_90) & (df["float_date"] >= end_date)]
col = "float_share" if "float_share" in df.columns else None
if col:
unlock_30 = safe_float(pd.to_numeric(df30[col], errors="coerce").sum(), 0.0) or 0.0
unlock_90 = safe_float(pd.to_numeric(df90[col], errors="coerce").sum(), 0.0) or 0.0
if float_shares and float_shares > 0:
# share_float 通常单位为万股,这里统一换算为股
result["unlock_30d_ratio"] = round(unlock_30 * 10000 / float_shares * 100.0, 4)
result["unlock_90d_ratio"] = round(unlock_90 * 10000 / float_shares * 100.0, 4)
except Exception as exc:
result["unlock_error"] = str(exc)
try:
df = pro().stk_holdertrade(ts_code=ts_code, start_date=back_30, end_date=end_date)
if df is not None and not df.empty:
de_col = "in_de" if "in_de" in df.columns else None
if de_col:
de_count = len(df[df[de_col] == "DE"])
result["reduction_density_30d"] = round(de_count / 30.0, 4)
else:
result["reduction_density_30d"] = round(len(df) / 30.0, 4)
except Exception as exc:
result["reduction_error"] = str(exc)
try:
df = pro().repurchase(ts_code=ts_code)
if df is not None and not df.empty:
date_col = "ann_date" if "ann_date" in df.columns else None
if date_col:
recent = df[df[date_col] >= back_90]
else:
recent = df
rep_ratio = 0.0
if total_shares and total_shares > 0 and "vol" in recent.columns:
vol = pd.to_numeric(recent["vol"], errors="coerce").fillna(0).sum()
rep_ratio = float(vol) * 10000 / total_shares * 100.0
else:
rep_ratio = len(recent) * 0.3
result["repurchase_ratio_90d"] = round(rep_ratio, 4)
except Exception as exc:
result["repurchase_error"] = str(exc)
return result
@retry_on_failure(max_retries=2, delay=1.0)
def get_index_constituents(index_name: str) -> list:
"""获取指数成分股"""
index_code = INDEX_CODE_MAP.get(index_name)
if not index_code:
return []
try:
end_date = datetime.now().strftime("%Y%m%d")
start_date = (datetime.now() - timedelta(days=90)).strftime("%Y%m%d")
df = pro().index_weight(index_code=index_code, start_date=start_date, end_date=end_date)
if df is None or df.empty:
return []
latest = str(df["trade_date"].max())
latest_df = df[df["trade_date"] == latest]
if latest_df.empty:
latest_df = df
codes = latest_df["con_code"].dropna().unique().tolist()
return [ts_to_symbol(c) for c in codes]
except Exception as exc:
print(f"获取指数成分股失败: {exc}")
return []
def get_all_a_stocks() -> list:
"""获取全部A股代码"""
try:
df = pro().stock_basic(exchange="", list_status="L", fields="symbol")
if df is not None and not df.empty:
return df["symbol"].dropna().unique().tolist()
return []
except Exception as exc:
print(f"获取全部A股失败: {exc}")
return []
def fetch_stock_data(
code: str,
data_type: str = "all",
years: int = 3,
use_cache: bool = True,
with_realtime: bool = False,
with_event_window: bool = False,
benchmark: str = "hs300",
realtime_window: int = 60,
event_window_pre: int = 1,
event_window_post: Sequence[int] = (1, 3, 5),
cache_ttl_min: int = 1440,
) -> dict:
"""获取单只股票的数据"""
code = normalize_symbol(code)
cache_key = data_type
if with_realtime:
cache_key += f"_rt_{benchmark}_{max(20, realtime_window)}"
if with_event_window:
post_key = "-".join(str(int(x)) for x in sorted(set(event_window_post)) if int(x) >= 1) or "1-3-5"
cache_key += f"_ew_{benchmark}_p{max(0, int(event_window_pre))}_w{post_key}"
if use_cache:
cached = load_cache(code, cache_key, ttl_minutes=cache_ttl_min)
if cached:
log(f"使用缓存数据: {code}")
return cached
result = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"data_type": data_type,
}
log(f"正在获取 {code} 的数据...")
if data_type in ["all", "basic", "news"]:
log(" - 获取基本信息...")
result["basic_info"] = get_stock_info(code)
if data_type in ["all", "financial"]:
log(" - 获取财务数据...")
result["financial_data"] = get_financial_data(code, years)
log(" - 获取财务指标...")
result["financial_indicators"] = get_financial_indicators(code)
log(" - 获取业绩与审计数据...")
result["performance_data"] = get_performance_data(code, years)
if data_type in ["all", "valuation"]:
log(" - 获取估值数据...")
result["valuation"] = get_valuation_data(code)
log(" - 获取价格数据...")
result["price"] = get_price_data(code, days=max(60, realtime_window))
if with_realtime:
if "basic_info" not in result:
log(" - 获取基本信息(实时模块依赖)...")
result["basic_info"] = get_stock_info(code)
if "price" not in result:
log(" - 获取价格数据(实时模块依赖)...")
result["price"] = get_price_data(code, days=max(60, realtime_window))
log(" - 获取资金流指标...")
result["flow_metrics"] = get_flow_metrics(code, days=20)
log(" - 获取筹码事件...")
result["chip_events"] = get_chip_events(code, basic_info=result.get("basic_info", {}))
log(f" - 计算实时指标(benchmark={benchmark})...")
benchmark_price = get_index_price_data(benchmark, days=max(60, realtime_window))
result["realtime_metrics"] = calculate_realtime_metrics(
price=result.get("price", {}),
benchmark_price=benchmark_price,
flow_metrics=result.get("flow_metrics", {}),
chip_events=result.get("chip_events", {}),
window=max(20, realtime_window),
)
result["realtime_metrics"]["benchmark"] = benchmark
if isinstance(benchmark_price, dict) and benchmark_price.get("error"):
result["realtime_metrics"]["benchmark_error"] = benchmark_price.get("error")
if with_event_window:
log(f" - 计算事件窗口(benchmark={benchmark})...")
result = attach_event_window(
result,
benchmark=benchmark,
pre_days=event_window_pre,
post_days=event_window_post,
)
if data_type in ["all", "holder"]:
log(" - 获取股东数据...")
result["holder"] = get_holder_data(code)
log(" - 获取分红数据...")
result["dividend"] = get_dividend_data(code)
if use_cache:
save_cache(code, cache_key, result)
log(f"数据获取完成: {code}")
return result
def fetch_multiple_stocks(
codes: list,
data_type: str = "basic",
with_realtime: bool = False,
with_event_window: bool = False,
benchmark: str = "hs300",
realtime_window: int = 60,
event_window_pre: int = 1,
event_window_post: Sequence[int] = (1, 3, 5),
cache_ttl_min: int = 1440,
) -> dict:
"""获取多只股票数据"""
result = {
"fetch_time": datetime.now().isoformat(),
"stocks": [],
"success_count": 0,
"fail_count": 0,
}
total = len(codes)
for i, code in enumerate(codes):
code = normalize_symbol(code)
log(f"[{i + 1}/{total}] 获取 {code}...")
try:
stock_data = fetch_stock_data(
code,
data_type,
use_cache=True,
with_realtime=with_realtime,
with_event_window=with_event_window,
benchmark=benchmark,
realtime_window=realtime_window,
event_window_pre=event_window_pre,
event_window_post=event_window_post,
cache_ttl_min=cache_ttl_min,
)
if "error" not in stock_data.get("basic_info", {}):
result["stocks"].append(stock_data)
result["success_count"] += 1
else:
result["fail_count"] += 1
except Exception as exc:
log(f" 获取失败: {exc}")
result["fail_count"] += 1
if i < total - 1:
time.sleep(0.2)
return result
def attach_news_data(
result: dict,
days: int = 7,
limit: int = 20,
news_sources: str = "",
news_provider: str = "auto",
brave_api_key: str = "",
tushare_token: str = "",
) -> dict:
"""为结果补充新闻与舆情数据。"""
code = result.get("code", "")
basic = result.get("basic_info", {})
name = basic.get("name", "")
try:
items = fetch_news(
code=code,
name=name,
days=days,
limit=limit,
provider=news_provider,
brave_api_key=brave_api_key,
tushare_token=tushare_token,
)
if news_sources:
allow_sources = {x.strip().lower() for x in news_sources.split(",") if x.strip()}
if allow_sources:
items = [x for x in items if x.get("source", "").lower() in allow_sources]
sentiment = analyze_news_sentiment(items)
result["news_items"] = items
result["news_sentiment"] = {
"analysis_time": sentiment.get("analysis_time"),
"news_count": sentiment.get("news_count"),
"overall_sentiment": sentiment.get("overall_sentiment"),
"risk_level": sentiment.get("risk_level"),
"risk_tag_count": sentiment.get("risk_tag_count"),
"top_negative_events": sentiment.get("top_negative_events"),
}
except Exception as exc:
result["news_items"] = []
result["news_sentiment"] = {
"analysis_time": datetime.now().isoformat(),
"news_count": 0,
"overall_sentiment": 0.0,
"risk_level": "低",
"risk_tag_count": {},
"top_negative_events": [],
"error": str(exc),
}
return result
def attach_event_window(
result: dict,
benchmark: str = "hs300",
pre_days: int = 1,
post_days: Sequence[int] = (1, 3, 5),
max_events: int = 40,
) -> dict:
"""为结果补充事件窗口分析。"""
if not isinstance(result, dict):
return result
code = result.get("code", "")
if "price" not in result or not result.get("price"):
if code:
log(" - 获取价格数据(事件窗口依赖)...")
result["price"] = get_price_data(code, days=max(90, max(post_days, default=5) * 12))
benchmark_price = get_index_price_data(benchmark, days=max(90, max(post_days, default=5) * 12))
candidates = collect_event_candidates(result, max_events=max_events * 2)
ew = calculate_event_window(
price=result.get("price", {}),
events=candidates,
benchmark_price=benchmark_price,
pre_days=pre_days,
post_days=post_days,
max_events=max_events,
)
ew["benchmark"] = benchmark
ew["event_sources"] = ["news_items", "performance_data.forecast", "performance_data.express", "performance_data.audit"]
if isinstance(benchmark_price, dict) and benchmark_price.get("error"):
ew["benchmark_error"] = benchmark_price.get("error")
result["event_window"] = ew
return result
def summarize_table(result: dict) -> str:
basic = result.get("basic_info", {})
price = result.get("price", {})
valuation = result.get("valuation", {})
sentiment = result.get("news_sentiment", {})
realtime = result.get("realtime_metrics", {})
event_window = result.get("event_window", {})
headers = ["代码", "名称", "行业", "最新价", "PE_TTM", "PB", "实时分", "事件窗分", "舆情", "数据类型"]
rows = [[
result.get("code", ""),
basic.get("name", ""),
basic.get("industry", ""),
price.get("latest_price", ""),
basic.get("pe_ttm", valuation.get("latest", {}).get("pe_ttm", "")),
basic.get("pb", valuation.get("latest", {}).get("pb", "")),
realtime.get("realtime_score", "-"),
event_window.get("event_window_score", "-"),
sentiment.get("risk_level", "-"),
result.get("data_type", ""),
]]
return format_table(headers, rows)
def main():
default_years = int(os.getenv("STOCK_ANALYSIS_DEFAULT_YEARS", "3"))
default_news_days = int(os.getenv("STOCK_ANALYSIS_DEFAULT_NEWS_DAYS", "7"))
default_news_limit = int(os.getenv("STOCK_ANALYSIS_DEFAULT_NEWS_LIMIT", "20"))
default_news_provider = os.getenv("STOCK_ANALYSIS_NEWS_PROVIDER", "auto")
default_cache_ttl = int(os.getenv("STOCK_ANALYSIS_CACHE_TTL_MIN", "1440"))
parser = argparse.ArgumentParser(description="A股数据获取工具")
parser.add_argument("--code", type=str, help="股票代码 (如: 600519)")
parser.add_argument("--codes", type=str, help="多个股票代码,逗号分隔 (如: 600519,000858)")
parser.add_argument(
"--data-type",
type=str,
default="basic",
choices=["all", "basic", "financial", "valuation", "holder", "news"],
help="数据类型 (默认: basic)",
)
parser.add_argument("--years", type=int, default=default_years, help=f"获取多少年的历史数据 (默认: {default_years})")
parser.add_argument("--scope", type=str, help="筛选范围: hs300/zz500/cyb/kcb/all")
parser.add_argument("--no-cache", action="store_true", help="不使用缓存")
parser.add_argument("--token", type=str, required=True, help="tushare token(必填)")
parser.add_argument("--with-news", action="store_true", help="附加最近新闻与舆情")
parser.add_argument("--news-days", type=int, default=default_news_days, help=f"新闻窗口天数 (默认: {default_news_days})")
parser.add_argument("--news-limit", type=int, default=default_news_limit, help=f"新闻最大条数 (默认: {default_news_limit})")
parser.add_argument("--news-sources", type=str, default="", help="新闻来源过滤,逗号分隔")
parser.add_argument("--news-provider", choices=["auto", "brave", "tushare", "rss"], default=default_news_provider, help=f"新闻源 (默认: {default_news_provider})")
parser.add_argument("--brave-api-key", type=str, default="", help="Brave Search API Key(news-provider=brave 时必填)")
parser.add_argument("--with-realtime", action="store_true", help="附加实时指标(趋势/确认/风险/筹码)")
parser.add_argument("--with-event-window", action="store_true", help="附加事件窗口分析(事件后1/3/5日反应)")
parser.add_argument("--benchmark", type=str, default="hs300", choices=["hs300", "zz500", "zz1000", "cyb", "kcb"], help="相对强弱基准指数")
parser.add_argument("--realtime-window", type=int, default=60, help="实时指标窗口(日)")
parser.add_argument("--event-window-pre", type=int, default=1, help="事件窗口前置天数(默认: 1)")
parser.add_argument("--event-window-post", type=str, default="1,3,5", help="事件窗口后验天数,逗号分隔(默认: 1,3,5)")
parser.add_argument("--cache-ttl-min", type=int, default=default_cache_ttl, help=f"缓存TTL分钟数 (默认: {default_cache_ttl})")
parser.add_argument("--format", choices=["json", "table"], default="json", help="输出格式")
parser.add_argument("--quiet", action="store_true", help="静默模式,仅输出结果")
parser.add_argument("--output", type=str, help="输出文件路径 (JSON)")
args = parser.parse_args()
if not str(args.token or "").strip():
parser.error("--token 不能为空(请确认 TUSHARE_TOKEN 已正确导出,或直接传入明文)")
global CLI_TOKEN, VERBOSE
CLI_TOKEN = args.token
VERBOSE = not args.quiet
event_window_post = parse_event_window_post_days(args.event_window_post)
result = {}
if args.code:
result = fetch_stock_data(
args.code,
args.data_type,
args.years,
use_cache=not args.no_cache,
with_realtime=args.with_realtime,
with_event_window=args.with_event_window,
benchmark=args.benchmark,
realtime_window=args.realtime_window,
event_window_pre=args.event_window_pre,
event_window_post=event_window_post,
cache_ttl_min=args.cache_ttl_min,
)
if args.data_type in ["all", "news"] or args.with_news:
log(" - 获取新闻与舆情...")
result = attach_news_data(
result,
days=args.news_days,
limit=args.news_limit,
news_sources=args.news_sources,
news_provider=args.news_provider,
brave_api_key=args.brave_api_key,
tushare_token=args.token,
)
if args.with_event_window:
log(" - 事件窗口重算(纳入新闻事件)...")
result = attach_event_window(
result,
benchmark=args.benchmark,
pre_days=args.event_window_pre,
post_days=event_window_post,
)
elif args.codes:
codes = [c.strip() for c in args.codes.split(",") if c.strip()]
result = fetch_multiple_stocks(
codes,
args.data_type,
with_realtime=args.with_realtime,
with_event_window=args.with_event_window,
benchmark=args.benchmark,
realtime_window=args.realtime_window,
event_window_pre=args.event_window_pre,
event_window_post=event_window_post,
cache_ttl_min=args.cache_ttl_min,
)
if args.with_news:
for item in result.get("stocks", []):
attach_news_data(
item,
days=args.news_days,
limit=args.news_limit,
news_sources=args.news_sources,
news_provider=args.news_provider,
brave_api_key=args.brave_api_key,
tushare_token=args.token,
)
if args.with_event_window:
attach_event_window(
item,
benchmark=args.benchmark,
pre_days=args.event_window_pre,
post_days=event_window_post,
)
elif args.scope:
if args.scope == "all":
codes = get_all_a_stocks()
else:
codes = get_index_constituents(args.scope)
result = {"scope": args.scope, "stocks": codes, "count": len(codes)}
else:
print("请提供 --code, --codes 或 --scope 参数")
sys.exit(1)
output = json.dumps(result, ensure_ascii=False, indent=2, default=str)
if args.output:
with open(args.output, "w", encoding="utf-8") as f:
f.write(output)
log(f"\n数据已保存到: {args.output}")
else:
if args.format == "table":
if args.scope and "stocks" in result:
rows = [[result.get("scope", ""), result.get("count", 0)]]
print(format_table(["范围", "股票数量"], rows))
elif args.codes and "stocks" in result:
rows = [[s.get("code", ""), s.get("basic_info", {}).get("name", "")] for s in result["stocks"]]
print(format_table(["代码", "名称"], rows))
else:
print(summarize_table(result))
else:
print(output)
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
环境变量加载工具:
- 优先读取当前进程环境变量
- 缺失时回退读取 ~/.aj-skills/.env
"""
from __future__ import annotations
import os
from pathlib import Path
from typing import Dict
AJ_ENV_PATH = Path.home() / ".aj-skills" / ".env"
def _parse_env_file(path: Path) -> Dict[str, str]:
if not path.exists() or not path.is_file():
return {}
data: Dict[str, str] = {}
for raw in path.read_text(encoding="utf-8").splitlines():
line = raw.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, value = line.split("=", 1)
key = key.strip()
value = value.strip().strip("'").strip('"')
if key:
data[key] = value
return data
def get_env(name: str, default: str = "") -> str:
val = os.getenv(name)
if val:
return val
from_file = _parse_env_file(AJ_ENV_PATH).get(name)
if from_file:
return from_file
return default
def get_tushare_token() -> str:
return get_env("TUSHARE_TOKEN") or get_env("TS_TOKEN")
def get_brave_api_key() -> str:
return get_env("BRAVE_API_KEY")
#!/usr/bin/env python3
"""
事件窗口分析模块(P1)
- 基于日线计算事件后 1/3/5 日价格反应
- 支持相对基准的超额收益(abnormal return)
- 事件源:新闻、业绩预告/快报、审计意见
"""
from __future__ import annotations
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Sequence, Tuple
def _safe_float(val, default: Optional[float] = None) -> Optional[float]:
if val is None:
return default
try:
if isinstance(val, str):
val = val.replace("%", "").replace(",", "").strip()
if not val:
return default
return float(val)
except Exception:
return default
def _clamp(value: float, lo: float = 0.0, hi: float = 100.0) -> float:
return max(lo, min(hi, value))
def _parse_datetime(value) -> Optional[datetime]:
if value is None:
return None
if isinstance(value, datetime):
return value.replace(tzinfo=None)
s = str(value).strip()
if not s:
return None
try:
return datetime.fromisoformat(s.replace("Z", "+00:00")).replace(tzinfo=None)
except Exception:
pass
digits = "".join(ch for ch in s if ch.isdigit())
if len(digits) >= 14:
try:
return datetime.strptime(digits[:14], "%Y%m%d%H%M%S")
except Exception:
pass
if len(digits) >= 8:
try:
return datetime.strptime(digits[:8], "%Y%m%d")
except Exception:
return None
return None
def _to_ymd(dt: datetime) -> str:
return dt.strftime("%Y%m%d")
def _event_anchor_date(event_dt: datetime) -> datetime:
# 对于含具体时间的事件,15:00后视为次一交易日生效
if event_dt.hour >= 15:
return (event_dt + timedelta(days=1)).replace(hour=0, minute=0, second=0, microsecond=0)
return event_dt.replace(hour=0, minute=0, second=0, microsecond=0)
def _extract_trade_rows(price: Dict) -> List[Dict]:
rows = price.get("price_data", []) if isinstance(price, dict) else []
out: List[Dict] = []
for row in rows:
d = _parse_datetime(row.get("日期"))
c = _safe_float(row.get("收盘"))
if d is None or c in [None, 0]:
continue
out.append({"date": d.replace(hour=0, minute=0, second=0, microsecond=0), "close": c, "raw": row})
out.sort(key=lambda x: x["date"])
return out
def _locate_trade_index(trade_dates: Sequence[datetime], anchor: datetime) -> Optional[int]:
for idx, d in enumerate(trade_dates):
if d >= anchor:
return idx
return None
def _pct_change(curr: Optional[float], prev: Optional[float]) -> Optional[float]:
if curr is None or prev in [None, 0]:
return None
return (curr / prev - 1.0) * 100.0
def _avg(values: List[Optional[float]]) -> Optional[float]:
arr = [x for x in values if x is not None]
if not arr:
return None
return sum(arr) / len(arr)
def _is_negative_audit(opinion: str) -> bool:
s = str(opinion or "")
return any(x in s for x in ["非标", "保留", "否定", "无法表示"])
def collect_event_candidates(stock_data: Dict, max_events: int = 60) -> List[Dict]:
"""从 stock_data 聚合事件候选集。"""
events: List[Dict] = []
news_items = stock_data.get("news_items", []) if isinstance(stock_data, dict) else []
for item in news_items or []:
event_dt = _parse_datetime(item.get("published_at") or item.get("date"))
if event_dt is None:
continue
title = str(item.get("title", "") or "").strip() or "新闻事件"
events.append(
{
"event_date": _to_ymd(event_dt),
"event_time": event_dt.isoformat(),
"event_type": "news",
"title": title,
"source": item.get("source", ""),
"sentiment": _safe_float(item.get("sentiment")),
"importance": _safe_float(item.get("importance"), 1.0) or 1.0,
}
)
perf = stock_data.get("performance_data", {}) if isinstance(stock_data, dict) else {}
forecast = perf.get("forecast", []) if isinstance(perf, dict) else []
for row in forecast[:12]:
event_dt = _parse_datetime(row.get("ann_date"))
if event_dt is None:
continue
pmin = _safe_float(row.get("p_change_min"))
pmax = _safe_float(row.get("p_change_max"))
avg_change = None
if pmin is not None and pmax is not None:
avg_change = (pmin + pmax) / 2.0
elif pmin is not None:
avg_change = pmin
elif pmax is not None:
avg_change = pmax
events.append(
{
"event_date": _to_ymd(event_dt),
"event_time": event_dt.isoformat(),
"event_type": "forecast",
"title": f"业绩预告:{row.get('type', '未披露类型')}(净利变动约 {avg_change if avg_change is not None else '-'}%)",
"source": "performance_data.forecast",
"sentiment": (avg_change or 0.0) / 50.0 if avg_change is not None else None,
"importance": 1.3,
}
)
express = perf.get("express", []) if isinstance(perf, dict) else []
for row in express[:12]:
event_dt = _parse_datetime(row.get("ann_date"))
if event_dt is None:
continue
yoy = _safe_float(row.get("yoy_net_profit"))
events.append(
{
"event_date": _to_ymd(event_dt),
"event_time": event_dt.isoformat(),
"event_type": "express",
"title": f"业绩快报:净利润同比 {yoy if yoy is not None else '-'}%",
"source": "performance_data.express",
"sentiment": (yoy or 0.0) / 50.0 if yoy is not None else None,
"importance": 1.2,
}
)
audit = perf.get("audit", []) if isinstance(perf, dict) else []
for row in audit[:8]:
event_dt = _parse_datetime(row.get("ann_date"))
if event_dt is None:
continue
opinion = str(row.get("audit_result", "") or row.get("audit_agency", "") or "审计披露")
events.append(
{
"event_date": _to_ymd(event_dt),
"event_time": event_dt.isoformat(),
"event_type": "audit",
"title": f"审计意见:{opinion}",
"source": "performance_data.audit",
"sentiment": -0.6 if _is_negative_audit(opinion) else 0.1,
"importance": 1.1,
}
)
# 去重:同日同标题只保留一条
dedup: Dict[Tuple[str, str], Dict] = {}
for item in events:
key = (str(item.get("event_date", "")), str(item.get("title", "")))
if key not in dedup:
dedup[key] = item
merged = list(dedup.values())
merged.sort(key=lambda x: x.get("event_date", ""), reverse=True)
return merged[:max_events]
def calculate_event_window(
price: Dict,
events: List[Dict],
benchmark_price: Optional[Dict] = None,
pre_days: int = 1,
post_days: Sequence[int] = (1, 3, 5),
max_events: int = 40,
) -> Dict:
"""计算事件窗口反应。"""
rows = _extract_trade_rows(price or {})
if not rows:
return {
"pre_days": max(0, int(pre_days)),
"post_days": [int(x) for x in post_days if int(x) >= 1],
"event_count": len(events or []),
"matched_event_count": 0,
"event_window_score": 50.0,
"events": [],
"error": "缺少有效价格序列",
}
bench_rows = _extract_trade_rows(benchmark_price or {})
bench_map = {x["date"]: x["close"] for x in bench_rows}
trade_dates = [x["date"] for x in rows]
closes = [x["close"] for x in rows]
post = sorted({int(x) for x in post_days if int(x) >= 1})
pre_days = max(0, int(pre_days))
analyzed: List[Dict] = []
for event in events[:max_events]:
event_dt = _parse_datetime(event.get("event_time") or event.get("event_date"))
if event_dt is None:
continue
anchor = _event_anchor_date(event_dt)
idx = _locate_trade_index(trade_dates, anchor)
if idx is None:
continue
base = closes[idx]
pre_idx = idx - pre_days
pre_ret = _pct_change(closes[idx - 1], closes[pre_idx]) if pre_days > 0 and pre_idx >= 0 and idx - 1 >= 0 else None
row = {
"event_date": event.get("event_date"),
"event_trade_date": _to_ymd(trade_dates[idx]),
"event_type": event.get("event_type"),
"title": event.get("title"),
"source": event.get("source"),
"sentiment": event.get("sentiment"),
"importance": event.get("importance"),
"pre_return_pct": round(pre_ret, 4) if pre_ret is not None else None,
}
event_date_obj = trade_dates[idx]
bench_base = bench_map.get(event_date_obj)
for n in post:
k = idx + n
if k >= len(closes):
row[f"post_{n}d_pct"] = None
row[f"abnormal_{n}d_pct"] = None
continue
post_ret = _pct_change(closes[k], base)
row[f"post_{n}d_pct"] = round(post_ret, 4) if post_ret is not None else None
bench_ret = None
if bench_base not in [None, 0]:
bench_px = bench_map.get(trade_dates[k])
bench_ret = _pct_change(bench_px, bench_base)
abn = (post_ret - bench_ret) if post_ret is not None and bench_ret is not None else None
row[f"abnormal_{n}d_pct"] = round(abn, 4) if abn is not None else None
core = row.get("abnormal_3d_pct")
if core is None:
core = row.get("post_3d_pct")
if core is None:
core = row.get("post_1d_pct")
score = 50.0 + (core or 0.0) * 4.0
sentiment = _safe_float(event.get("sentiment"))
if sentiment is not None:
score += max(-10.0, min(10.0, sentiment * 10.0))
row["event_impact_score"] = round(_clamp(score), 2)
analyzed.append(row)
analyzed.sort(key=lambda x: x.get("event_trade_date", ""), reverse=True)
def collect_metric(metric_key: str) -> List[Optional[float]]:
return [_safe_float(x.get(metric_key)) for x in analyzed]
avg_post_1 = _avg(collect_metric("post_1d_pct"))
avg_post_3 = _avg(collect_metric("post_3d_pct"))
avg_post_5 = _avg(collect_metric("post_5d_pct"))
avg_abn_1 = _avg(collect_metric("abnormal_1d_pct"))
avg_abn_3 = _avg(collect_metric("abnormal_3d_pct"))
avg_abn_5 = _avg(collect_metric("abnormal_5d_pct"))
post_3_values = [x for x in collect_metric("post_3d_pct") if x is not None]
positive_ratio_3d = (len([x for x in post_3_values if x > 0]) / len(post_3_values)) if post_3_values else None
post_5_values = [x for x in collect_metric("post_5d_pct") if x is not None]
worst_5d = min(post_5_values) if post_5_values else None
ev_score = 50.0
if avg_abn_3 is not None:
ev_score += avg_abn_3 * 6.0
elif avg_post_3 is not None:
ev_score += avg_post_3 * 4.0
if positive_ratio_3d is not None:
ev_score += (positive_ratio_3d - 0.5) * 40.0
if worst_5d is not None and worst_5d < -5:
ev_score -= min(20.0, abs(worst_5d + 5.0) * 1.8)
ranked = []
for item in analyzed:
key = _safe_float(item.get("abnormal_3d_pct"))
if key is None:
key = _safe_float(item.get("post_3d_pct"))
if key is None:
key = _safe_float(item.get("post_1d_pct"), 0.0)
ranked.append((key or 0.0, item))
ranked.sort(key=lambda x: x[0], reverse=True)
top_positive = [x[1] for x in ranked[:5] if x[0] > 0]
top_negative = [x[1] for x in sorted(ranked, key=lambda x: x[0])[:5] if x[0] < 0]
summary = {
"avg_post_1d_pct": round(avg_post_1, 4) if avg_post_1 is not None else None,
"avg_post_3d_pct": round(avg_post_3, 4) if avg_post_3 is not None else None,
"avg_post_5d_pct": round(avg_post_5, 4) if avg_post_5 is not None else None,
"avg_abnormal_1d_pct": round(avg_abn_1, 4) if avg_abn_1 is not None else None,
"avg_abnormal_3d_pct": round(avg_abn_3, 4) if avg_abn_3 is not None else None,
"avg_abnormal_5d_pct": round(avg_abn_5, 4) if avg_abn_5 is not None else None,
"positive_ratio_3d": round(positive_ratio_3d, 4) if positive_ratio_3d is not None else None,
"worst_post_5d_pct": round(worst_5d, 4) if worst_5d is not None else None,
}
return {
"pre_days": pre_days,
"post_days": post,
"benchmark": (benchmark_price or {}).get("benchmark"),
"event_count": len(events or []),
"matched_event_count": len(analyzed),
"coverage_ratio": round(len(analyzed) / len(events), 4) if events else 0.0,
"summary": summary,
"event_window_score": round(_clamp(ev_score), 2),
"top_positive_events": top_positive,
"top_negative_events": top_negative,
"events": analyzed,
}
#!/usr/bin/env python3
"""分模块抓取:基础信息。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
def main():
parser = argparse.ArgumentParser(description="抓取基础信息 basic.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "basic",
"basic_info": df.get_stock_info(code),
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:事件窗口分析。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
from event_window import collect_event_candidates, calculate_event_window
from news_fetcher import fetch_news
def _parse_post_days(raw: str):
values = []
for p in str(raw or "").split(","):
p = p.strip()
if not p:
continue
try:
v = int(p)
if v >= 1:
values.append(v)
except Exception:
continue
if not values:
return (1, 3, 5)
return tuple(sorted(set(values)))
def main():
parser = argparse.ArgumentParser(description="抓取事件窗口模块 event_window.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--name", default="", help="股票名称,可选")
parser.add_argument("--years", type=int, default=2, help="业绩事件回溯年数")
parser.add_argument("--days", type=int, default=7, help="新闻窗口天数")
parser.add_argument("--limit", type=int, default=20, help="新闻最大条数")
parser.add_argument("--provider", choices=["auto", "brave", "tushare", "rss"], default="auto", help="新闻源")
parser.add_argument("--brave-api-key", default="", help="Brave API Key(provider=brave 时必填)")
parser.add_argument("--benchmark", default="hs300", choices=["hs300", "zz500", "zz1000", "cyb", "kcb"], help="基准指数")
parser.add_argument("--pre-days", type=int, default=1, help="前置窗口天数")
parser.add_argument("--post-days", default="1,3,5", help="后验窗口天数,逗号分隔")
parser.add_argument("--max-events", type=int, default=40, help="最大事件数")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
if not str(args.token or "").strip():
parser.error("--token 不能为空(请确认 TUSHARE_TOKEN 已正确导出,或直接传入明文)")
if args.provider == "tushare" and not str(args.token or "").strip():
parser.error("provider=tushare 时,--token 不能为空")
if args.provider == "brave" and not str(args.brave_api_key or "").strip():
parser.error("provider=brave 时,--brave-api-key 不能为空")
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
post_days = _parse_post_days(args.post_days)
basic_info = df.get_stock_info(code)
name = args.name or basic_info.get("name", "")
price = df.get_price_data(code, days=max(120, max(post_days) * 20))
benchmark_price = df.get_index_price_data(args.benchmark, days=max(120, max(post_days) * 20))
performance_data = df.get_performance_data(code, years=max(1, args.years))
news_items = fetch_news(
code=code,
name=name,
days=max(1, args.days),
limit=max(1, args.limit),
provider=args.provider,
brave_api_key=args.brave_api_key,
tushare_token=args.token,
)
event_source = {
"code": code,
"news_items": news_items,
"performance_data": performance_data,
}
candidates = collect_event_candidates(event_source, max_events=max(10, args.max_events * 2))
event_window = calculate_event_window(
price=price,
events=candidates,
benchmark_price=benchmark_price,
pre_days=max(0, args.pre_days),
post_days=post_days,
max_events=max(10, args.max_events),
)
event_window["benchmark"] = args.benchmark
if isinstance(benchmark_price, dict) and benchmark_price.get("error"):
event_window["benchmark_error"] = benchmark_price.get("error")
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "event_window",
"event_window": event_window,
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:财务数据(报表/指标/业绩与审计)。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
def main():
parser = argparse.ArgumentParser(description="抓取财务模块 financial.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--years", type=int, default=3, help="历史年数")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "financial",
"financial_data": df.get_financial_data(code, years=max(1, args.years)),
"financial_indicators": df.get_financial_indicators(code),
"performance_data": df.get_performance_data(code, years=max(1, args.years)),
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:股东与分红数据。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
def main():
parser = argparse.ArgumentParser(description="抓取股东模块 holder.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "holder",
"holder": df.get_holder_data(code),
"dividend": df.get_dividend_data(code),
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:新闻与舆情数据。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
from sentiment_analyzer import analyze_news_sentiment
from news_fetcher import fetch_news
def main():
parser = argparse.ArgumentParser(description="抓取新闻与舆情模块 news.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--name", default="", help="股票名称,可选")
parser.add_argument("--days", type=int, default=7, help="新闻窗口天数")
parser.add_argument("--limit", type=int, default=20, help="新闻最大条数")
parser.add_argument("--provider", choices=["auto", "brave", "tushare", "rss"], default="auto", help="新闻源")
parser.add_argument("--token", default="", help="tushare token(provider=tushare 时必填)")
parser.add_argument("--brave-api-key", default="", help="Brave API Key(provider=brave 时必填)")
parser.add_argument("--news-sources", default="", help="来源过滤,逗号分隔")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
if args.provider == "tushare" and not str(args.token or "").strip():
parser.error("provider=tushare 时,--token 不能为空")
if args.provider == "brave" and not str(args.brave_api_key or "").strip():
parser.error("provider=brave 时,--brave-api-key 不能为空")
items = fetch_news(
code=args.code,
name=args.name,
days=max(1, args.days),
limit=max(1, args.limit),
provider=args.provider,
brave_api_key=args.brave_api_key,
tushare_token=args.token,
)
if args.news_sources:
allow_sources = {x.strip().lower() for x in args.news_sources.split(",") if x.strip()}
if allow_sources:
items = [x for x in items if x.get("source", "").lower() in allow_sources]
sentiment = analyze_news_sentiment(items)
payload = {
"code": args.code,
"fetch_time": datetime.now().isoformat(),
"section": "news",
"news_items": items,
"news_sentiment": {
"analysis_time": sentiment.get("analysis_time"),
"news_count": sentiment.get("news_count"),
"overall_sentiment": sentiment.get("overall_sentiment"),
"risk_level": sentiment.get("risk_level"),
"risk_tag_count": sentiment.get("risk_tag_count"),
"top_negative_events": sentiment.get("top_negative_events"),
},
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:价格行情数据。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
def main():
parser = argparse.ArgumentParser(description="抓取价格模块 price.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--days", type=int, default=120, help="抓取天数")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "price",
"price": df.get_price_data(code, days=max(60, args.days)),
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:实时指标(趋势/确认/风险/筹码)。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
from realtime_metrics import calculate_realtime_metrics
def main():
parser = argparse.ArgumentParser(description="抓取实时模块 realtime.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--benchmark", default="hs300", choices=["hs300", "zz500", "zz1000", "cyb", "kcb"], help="相对强弱基准指数")
parser.add_argument("--window", type=int, default=60, help="窗口(日)")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
window = max(20, int(args.window))
basic_info = df.get_stock_info(code)
price = df.get_price_data(code, days=max(60, window))
flow_metrics = df.get_flow_metrics(code, days=20)
chip_events = df.get_chip_events(code, basic_info=basic_info)
benchmark_price = df.get_index_price_data(args.benchmark, days=max(60, window))
realtime_metrics = calculate_realtime_metrics(
price=price,
benchmark_price=benchmark_price,
flow_metrics=flow_metrics,
chip_events=chip_events,
window=window,
)
realtime_metrics["benchmark"] = args.benchmark
if isinstance(benchmark_price, dict) and benchmark_price.get("error"):
realtime_metrics["benchmark_error"] = benchmark_price.get("error")
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "realtime",
"flow_metrics": flow_metrics,
"chip_events": chip_events,
"realtime_metrics": realtime_metrics,
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""分模块抓取:估值数据。"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from pathlib import Path
def main():
parser = argparse.ArgumentParser(description="抓取估值模块 valuation.json")
parser.add_argument("--code", required=True, help="股票代码,如 600519")
parser.add_argument("--token", required=True, help="tushare token(必填)")
parser.add_argument("--output", required=True, help="输出路径")
parser.add_argument("--quiet", action="store_true", help="静默模式")
args = parser.parse_args()
import data_fetcher as df
df.CLI_TOKEN = args.token
df.VERBOSE = not args.quiet
code = df.normalize_symbol(args.code)
payload = {
"code": code,
"fetch_time": datetime.now().isoformat(),
"section": "valuation",
"valuation": df.get_valuation_data(code),
}
out = Path(args.output)
out.parent.mkdir(parents=True, exist_ok=True)
out.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=str), encoding="utf-8")
if not args.quiet:
print(f"已保存: {out}")
if __name__ == "__main__":
main()
#!/usr/bin/env python3
"""
新闻舆情分析模块(MVP)
- 词典规则法进行情绪与风险标签识别
"""
from __future__ import annotations
import argparse
import json
from datetime import datetime
from typing import Dict, List
POSITIVE_WORDS = [
"增长", "创新高", "突破", "上调", "增持", "回购", "中标", "利好", "超预期", "改善",
]
NEGATIVE_WORDS = [
"下滑", "亏损", "暴跌", "诉讼", "减持", "处罚", "调查", "违约", "停产", "利空", "风险",
]
RISK_PATTERNS = {
"监管": ["处罚", "调查", "问询", "立案", "监管"],
"诉讼": ["诉讼", "仲裁", "索赔"],
"经营": ["亏损", "下滑", "违约", "停产", "裁员"],
"股东行为": ["减持", "质押", "冻结"],
}
def _score_text(text: str) -> float:
if not text:
return 0.0
pos = sum(1 for w in POSITIVE_WORDS if w in text)
neg = sum(1 for w in NEGATIVE_WORDS if w in text)
raw = pos - neg
if raw == 0:
return 0.0
# 压缩到 [-1, 1]
return max(-1.0, min(1.0, raw / 5.0))
def _extract_risk_tags(text: str) -> List[str]:
tags = []
for tag, patterns in RISK_PATTERNS.items():
if any(p in text for p in patterns):
tags.append(tag)
return tags
def analyze_news_sentiment(news_items: List[Dict]) -> Dict:
scored_items = []
all_scores: List[float] = []
tag_count: Dict[str, int] = {}
for item in news_items or []:
title = item.get("title", "")
summary = item.get("summary", "")
text = f"{title} {summary}".strip()
sentiment = _score_text(text)
tags = _extract_risk_tags(text)
for tag in tags:
tag_count[tag] = tag_count.get(tag, 0) + 1
enriched = {
**item,
"sentiment": sentiment,
"risk_tags": tags,
}
scored_items.append(enriched)
all_scores.append(sentiment)
overall = sum(all_scores) / len(all_scores) if all_scores else 0.0
if overall <= -0.3:
level = "高"
elif overall <= -0.1:
level = "中"
else:
level = "低"
top_negative = sorted(
[x for x in scored_items if x.get("sentiment", 0) < 0],
key=lambda x: x.get("sentiment", 0),
)[:5]
return {
"analysis_time": datetime.now().isoformat(),
"news_count": len(scored_items),
"overall_sentiment": round(overall, 4),
"risk_level": level,
"risk_tag_count": tag_count,
"top_negative_events": top_negative,
"news_items": scored_items,
}
def main():
parser = argparse.ArgumentParser(description="新闻舆情分析工具")
parser.add_argument("--input", required=True, help="输入新闻 JSON 文件")
parser.add_argument("--output", help="输出结果 JSON 文件")
args = parser.parse_args()
with open(args.input, "r", encoding="utf-8") as f:
data = json.load(f)
news_items = data.get("news_items", data if isinstance(data, list) else [])
result = analyze_news_sentiment(news_items)
output = json.dumps(result, ensure_ascii=False, indent=2)
if args.output:
with open(args.output, "w", encoding="utf-8") as f:
f.write(output)
else:
print(output)
if __name__ == "__main__":
main()
import json
from pathlib import Path
import sys
ROOT = Path(__file__).resolve().parents[1]
SCRIPTS = ROOT / "scripts"
sys.path.insert(0, str(SCRIPTS))
from data_contract import validate_stock_data
def load_fixture():
return json.loads((Path(__file__).parent / "fixtures" / "sample_stock_data.json").read_text(encoding="utf-8"))
def test_validate_stock_data_success():
ok, errors = validate_stock_data(load_fixture(), required_sections=["financial_data", "price", "valuation"])
assert ok
assert errors == []
def test_validate_stock_data_missing_key():
data = load_fixture()
data.pop("basic_info")
ok, errors = validate_stock_data(data)
assert not ok
assert any("basic_info" in err for err in errors)