"""Scan local A-share K-line/profile data for market reverse-gap candidates.
|
|
The script uses local MySQL read-only tables:
|
|
- a_share_profile_snapshot
|
- a_share_daily_price
|
|
It writes market manifestation audit artifacts into the industry case layout.
|
Use the result as a gap-finding input, not as an investment conclusion.
|
"""
|
|
from __future__ import annotations
|
|
import argparse
|
import csv
|
import hashlib
|
import json
|
import os
|
import re
|
import subprocess
|
import sys
|
from datetime import datetime, timezone
|
from pathlib import Path
|
|
|
SCRIPT_PATH = Path(__file__).resolve()
|
PROJECT_ROOT = SCRIPT_PATH.parents[2]
|
DEFAULT_MYSQL_EXE = Path(os.environ.get("TIANXIA_MYSQL_EXE", "M:/mysql/server/bin/mysql.exe"))
|
|
|
def now_iso() -> str:
|
return datetime.now(timezone.utc).astimezone().isoformat(timespec="seconds")
|
|
|
def safe_name(value: str, fallback: str = "artifact", max_len: int = 96) -> str:
|
value = re.sub(r"[^A-Za-z0-9._-]+", "_", value.strip())
|
value = value.strip("._-")
|
return (value or fallback)[:max_len]
|
|
|
def sql_literal(value: object) -> str:
|
if value is None:
|
return "NULL"
|
text = str(value)
|
return "'" + text.replace("\\", "\\\\").replace("'", "''") + "'"
|
|
|
def sha256_text(value: str) -> str:
|
return hashlib.sha256(value.encode("utf-8")).hexdigest()
|
|
|
def project_relative(path: Path) -> str:
|
try:
|
return path.resolve().relative_to(PROJECT_ROOT.resolve()).as_posix()
|
except ValueError:
|
return path.resolve().as_posix()
|
|
|
def keyword_predicate(keywords: list[str]) -> str:
|
clauses = []
|
for keyword in keywords:
|
pattern = "%" + keyword.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%"
|
lit = sql_literal(pattern)
|
clauses.append(
|
"("
|
f"p.industry_l1 LIKE {lit} ESCAPE '\\\\' OR "
|
f"p.industry_l2 LIKE {lit} ESCAPE '\\\\' OR "
|
f"p.concepts_raw LIKE {lit} ESCAPE '\\\\' OR "
|
f"p.main_business_raw LIKE {lit} ESCAPE '\\\\' OR "
|
f"p.stock_name LIKE {lit} ESCAPE '\\\\'"
|
")"
|
)
|
return " OR ".join(clauses)
|
|
|
def build_sql(args: argparse.Namespace) -> str:
|
if not args.keyword:
|
raise SystemExit("pass at least one --keyword")
|
pred = keyword_predicate(args.keyword)
|
start = sql_literal(args.start_date)
|
end = sql_literal(args.end_date)
|
min_return = float(args.min_return_pct)
|
strong_day_pct = float(args.strong_day_pct)
|
amount_ratio = float(args.amount_ratio)
|
return f"""
|
WITH latest_profile AS (
|
SELECT MAX(snapshot_date) AS snapshot_date
|
FROM a_share_profile_snapshot
|
WHERE snapshot_date <= {end}
|
),
|
universe AS (
|
SELECT
|
p.symbol,
|
p.stock_name,
|
p.industry_l1,
|
p.industry_l2,
|
p.concepts_raw,
|
p.main_business_raw
|
FROM a_share_profile_snapshot p
|
JOIN latest_profile lp ON lp.snapshot_date = p.snapshot_date
|
WHERE p.status = 'OK'
|
AND ({pred})
|
),
|
price_with_prev AS (
|
SELECT
|
d.symbol,
|
d.trade_date,
|
d.open_price,
|
d.high_price,
|
d.low_price,
|
d.close_price,
|
d.volume,
|
d.amount,
|
LAG(d.close_price) OVER (PARTITION BY d.symbol ORDER BY d.trade_date) AS prev_close,
|
AVG(d.amount) OVER (
|
PARTITION BY d.symbol
|
ORDER BY d.trade_date
|
ROWS BETWEEN 20 PRECEDING AND 1 PRECEDING
|
) AS avg_amount_20
|
FROM a_share_daily_price d
|
JOIN universe u ON u.symbol = d.symbol
|
WHERE d.trade_date BETWEEN DATE_SUB({start}, INTERVAL 30 DAY) AND {end}
|
),
|
window_rows AS (
|
SELECT
|
u.stock_name,
|
u.industry_l1,
|
u.industry_l2,
|
u.concepts_raw,
|
u.main_business_raw,
|
p.*,
|
CASE
|
WHEN p.prev_close IS NULL OR p.prev_close = 0 THEN NULL
|
ELSE (p.close_price / p.prev_close - 1) * 100
|
END AS daily_return_pct,
|
CASE
|
WHEN p.avg_amount_20 IS NULL OR p.avg_amount_20 = 0 THEN NULL
|
ELSE p.amount / p.avg_amount_20
|
END AS amount_ratio_20
|
FROM price_with_prev p
|
JOIN universe u ON u.symbol = p.symbol
|
WHERE p.trade_date BETWEEN {start} AND {end}
|
),
|
ranked AS (
|
SELECT
|
w.*,
|
ROW_NUMBER() OVER (PARTITION BY w.symbol ORDER BY w.trade_date) AS rn_first,
|
ROW_NUMBER() OVER (PARTITION BY w.symbol ORDER BY w.trade_date DESC) AS rn_last
|
FROM window_rows w
|
),
|
agg AS (
|
SELECT
|
symbol,
|
MAX(stock_name) AS company_name,
|
MAX(industry_l1) AS industry_l1,
|
MAX(industry_l2) AS industry_l2,
|
MAX(concepts_raw) AS concepts_raw,
|
MAX(main_business_raw) AS main_business_raw,
|
MIN(trade_date) AS window_start,
|
MAX(trade_date) AS window_end,
|
MAX(CASE WHEN rn_first = 1 THEN close_price END) AS first_close,
|
MAX(CASE WHEN rn_last = 1 THEN close_price END) AS last_close,
|
COUNT(*) AS trading_days,
|
SUM(CASE WHEN daily_return_pct >= {strong_day_pct} THEN 1 ELSE 0 END) AS strong_up_days,
|
MAX(daily_return_pct) AS max_daily_return_pct,
|
MAX(amount_ratio_20) AS max_amount_ratio_20,
|
MAX(amount) AS max_amount
|
FROM ranked
|
GROUP BY symbol
|
),
|
scored AS (
|
SELECT
|
*,
|
CASE
|
WHEN first_close IS NULL OR first_close = 0 THEN NULL
|
ELSE (last_close / first_close - 1) * 100
|
END AS period_return_pct
|
FROM agg
|
)
|
SELECT
|
{sql_literal(args.scan_id)} AS scan_id,
|
{sql_literal(args.industry_case)} AS industry_case,
|
symbol,
|
company_name,
|
industry_l1,
|
industry_l2,
|
window_start,
|
window_end,
|
trading_days,
|
ROUND(first_close, 4) AS first_close,
|
ROUND(last_close, 4) AS last_close,
|
ROUND(period_return_pct, 4) AS period_return_pct,
|
strong_up_days,
|
ROUND(max_daily_return_pct, 4) AS max_daily_return_pct,
|
ROUND(max_amount_ratio_20, 4) AS max_amount_ratio_20,
|
CASE
|
WHEN period_return_pct >= {min_return} OR strong_up_days >= 1 OR max_amount_ratio_20 >= {amount_ratio}
|
THEN 'STRONG_MANIFESTATION'
|
ELSE 'WEAK_OR_NORMAL'
|
END AS manifestation_type,
|
'REVIEW_REQUIRED' AS scope_type,
|
'REVIEW_REQUIRED' AS suspected_missing_track,
|
'MARKET_REVERSE_GAP_SCAN' AS gap_type,
|
CASE
|
WHEN period_return_pct >= {min_return} OR strong_up_days >= 1 OR max_amount_ratio_20 >= {amount_ratio}
|
THEN 1
|
ELSE 0
|
END AS supplement_required_flag,
|
'Review core business against industry scheme before adding to evidence queue' AS supplement_question,
|
'DRAFT_FOR_REVIEW' AS review_status,
|
LEFT(concepts_raw, 500) AS concepts_excerpt,
|
LEFT(main_business_raw, 500) AS main_business_excerpt
|
FROM scored
|
WHERE period_return_pct >= {min_return}
|
OR strong_up_days >= 1
|
OR max_amount_ratio_20 >= {amount_ratio}
|
ORDER BY supplement_required_flag DESC, period_return_pct DESC, strong_up_days DESC, max_amount_ratio_20 DESC;
|
""".strip()
|
|
|
def run_mysql(args: argparse.Namespace, sql: str) -> str:
|
password = os.environ.get("TIANXIA_MYSQL_PASSWORD")
|
if not password:
|
raise SystemExit("TIANXIA_MYSQL_PASSWORD is required in the environment")
|
mysql_exe = Path(args.mysql_exe)
|
if not mysql_exe.exists():
|
raise SystemExit(f"mysql executable not found: {mysql_exe}")
|
cmd = [
|
str(mysql_exe),
|
"--protocol=TCP",
|
"--batch",
|
"--raw",
|
"--quick",
|
"--default-character-set=utf8mb4",
|
f"--host={args.host}",
|
f"--port={args.port}",
|
f"--user={args.user}",
|
args.database,
|
"-e",
|
sql,
|
]
|
env = os.environ.copy()
|
env["MYSQL_PWD"] = password
|
proc = subprocess.run(cmd, text=True, capture_output=True, env=env, check=False)
|
if proc.returncode != 0:
|
raise SystemExit(proc.stderr.strip() or proc.stdout.strip())
|
return proc.stdout
|
|
|
def write_csv_from_tsv(tsv: str, out_path: Path) -> int:
|
out_path.parent.mkdir(parents=True, exist_ok=True)
|
row_count = 0
|
with out_path.open("w", encoding="utf-8-sig", newline="") as fh:
|
writer = csv.writer(fh)
|
for line_index, line in enumerate(tsv.splitlines()):
|
writer.writerow(line.split("\t"))
|
if line_index > 0:
|
row_count += 1
|
return row_count
|
|
|
def append_manifest(path: Path, row: dict[str, str]) -> None:
|
exists = path.exists()
|
path.parent.mkdir(parents=True, exist_ok=True)
|
with path.open("a", encoding="utf-8-sig", newline="") as fh:
|
writer = csv.DictWriter(fh, fieldnames=list(row.keys()))
|
if not exists:
|
writer.writeheader()
|
writer.writerow(row)
|
|
|
def scan(args: argparse.Namespace) -> int:
|
if not args.scan_id:
|
args.scan_id = f"{safe_name(args.run_id, 'run')}_MARKET_GAP_SCAN"
|
industry_root = PROJECT_ROOT / "ana-data" / "cases" / args.industry_case
|
supplement_dir = industry_root / "supplement"
|
manifest_dir = industry_root / "manifest"
|
supplement_dir.mkdir(parents=True, exist_ok=True)
|
manifest_dir.mkdir(parents=True, exist_ok=True)
|
|
run_safe = safe_name(args.run_id, "run")
|
sql = build_sql(args)
|
sql_hash = sha256_text(sql)
|
audit_path = supplement_dir / f"market_manifestation_gap_audit_{run_safe}_{sql_hash[:12]}.csv"
|
summary_path = supplement_dir / f"market_manifestation_gap_summary_{run_safe}_{sql_hash[:12]}.json"
|
sql_path = manifest_dir / f"market_manifestation_gap_query_{run_safe}_{sql_hash[:12]}.sql"
|
manifest_path = manifest_dir / f"market_manifestation_gap_manifest_{run_safe}.csv"
|
|
tsv = run_mysql(args, sql)
|
row_count = write_csv_from_tsv(tsv, audit_path)
|
sql_path.write_text(sql + "\n", encoding="utf-8")
|
summary = {
|
"scan_id": args.scan_id,
|
"case_id": args.case_id,
|
"batch_id": args.batch_id,
|
"run_id": args.run_id,
|
"industry_case": args.industry_case,
|
"start_date": args.start_date,
|
"end_date": args.end_date,
|
"keywords": args.keyword,
|
"min_return_pct": args.min_return_pct,
|
"strong_day_pct": args.strong_day_pct,
|
"amount_ratio": args.amount_ratio,
|
"row_count": row_count,
|
"audit_path": project_relative(audit_path),
|
"sql_sha256": sql_hash,
|
"generated_at": now_iso(),
|
"readout": "GAP_SCAN_ONLY_NOT_INVESTMENT_CONCLUSION",
|
}
|
summary_path.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8")
|
append_manifest(
|
manifest_path,
|
{
|
"scan_id": args.scan_id,
|
"case_id": args.case_id,
|
"batch_id": args.batch_id,
|
"run_id": args.run_id,
|
"industry_case": args.industry_case,
|
"start_date": args.start_date,
|
"end_date": args.end_date,
|
"keywords": "|".join(args.keyword),
|
"sql_sha256": sql_hash,
|
"sql_relative_path": project_relative(sql_path),
|
"audit_relative_path": project_relative(audit_path),
|
"summary_relative_path": project_relative(summary_path),
|
"row_count": str(row_count),
|
"generated_at": summary["generated_at"],
|
"review_status": "DRAFT_FOR_REVIEW",
|
},
|
)
|
print(f"audit={project_relative(audit_path)} rows={row_count}")
|
print(f"summary={project_relative(summary_path)}")
|
return 0
|
|
|
def parse_args(argv: list[str]) -> argparse.Namespace:
|
parser = argparse.ArgumentParser(description=__doc__)
|
parser.add_argument("--industry-case", required=True)
|
parser.add_argument("--case-id", required=True)
|
parser.add_argument("--batch-id", required=True)
|
parser.add_argument("--run-id", required=True)
|
parser.add_argument("--scan-id")
|
parser.add_argument("--start-date", required=True)
|
parser.add_argument("--end-date", required=True)
|
parser.add_argument("--keyword", action="append", required=True)
|
parser.add_argument("--min-return-pct", type=float, default=20.0)
|
parser.add_argument("--strong-day-pct", type=float, default=9.5)
|
parser.add_argument("--amount-ratio", type=float, default=2.0)
|
parser.add_argument("--mysql-exe", default=str(DEFAULT_MYSQL_EXE))
|
parser.add_argument("--host", default=os.environ.get("TIANXIA_MYSQL_HOST", "127.0.0.1"))
|
parser.add_argument("--port", default=os.environ.get("TIANXIA_MYSQL_PORT", "3306"))
|
parser.add_argument("--user", default=os.environ.get("TIANXIA_MYSQL_USER", "root"))
|
parser.add_argument("--database", default=os.environ.get("TIANXIA_MYSQL_DB", "tianxia"))
|
return parser.parse_args(argv)
|
|
|
def main(argv: list[str] | None = None) -> int:
|
args = parse_args(argv or sys.argv[1:])
|
return scan(args)
|
|
|
if __name__ == "__main__":
|
raise SystemExit(main())
|