63 lines
3.3 KiB
Python
63 lines
3.3 KiB
Python
|
|
"""Collect one real sector into an isolated PostgreSQL snapshot store.
|
||
|
|
|
||
|
|
Run with the server's uv environment from zhixing-server. No production
|
||
|
|
publication is created; facts must be read back from this store for validation.
|
||
|
|
"""
|
||
|
|
from datetime import date
|
||
|
|
from pathlib import Path
|
||
|
|
from urllib.parse import urlsplit, urlunsplit
|
||
|
|
import json
|
||
|
|
import time
|
||
|
|
import tushare as ts
|
||
|
|
from zhixing_server.bootstrap.config import Settings
|
||
|
|
from zhixing_server.modules.sector_radar.domain.source import build_source_snapshot
|
||
|
|
from zhixing_server.modules.sector_radar.infrastructure.postgres import PostgresSectorRadarRepository
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
"""Persist provider responses before inspecting their values; redact errors."""
|
||
|
|
settings = Settings(_env_file='../.env')
|
||
|
|
url = urlsplit(settings.database_url)
|
||
|
|
host = url.netloc.rsplit('@', 1)[0] + '@127.0.0.1:5433'
|
||
|
|
database = urlunsplit((url.scheme, host, '/radar_detail_selftest_0906', url.query, ''))
|
||
|
|
repository = PostgresSectorRadarRepository(database, max_connections=2)
|
||
|
|
client = ts.pro_api(settings.tushare_token)
|
||
|
|
target = date(2026, 9, 4)
|
||
|
|
counts = []
|
||
|
|
|
||
|
|
def collect(api, fields, **params):
|
||
|
|
"""Store each raw response atomically and return persisted rows."""
|
||
|
|
time.sleep(0.25)
|
||
|
|
frame = client.query(api, fields=fields, **params)
|
||
|
|
snapshot = build_source_snapshot(api_name=api, params=params,
|
||
|
|
rows=frame.to_dict('records'), target_trade_date=target)
|
||
|
|
repository.save_source_snapshots((snapshot,))
|
||
|
|
counts.append({'api':api, 'rows':snapshot.row_count, 'snapshot':snapshot.snapshot_id})
|
||
|
|
return snapshot.rows
|
||
|
|
|
||
|
|
try:
|
||
|
|
collect('dc_index','ts_code,trade_date,name,idx_type,level,pct_change,leading_code',
|
||
|
|
trade_date='20260904', ts_code='BK1147.DC')
|
||
|
|
members = collect('dc_member','trade_date,ts_code,con_code,name',
|
||
|
|
trade_date='20260904', ts_code='BK1147.DC')
|
||
|
|
collect('stock_basic','ts_code,symbol,name,market,exchange,list_status,list_date,delist_date',list_status='L')
|
||
|
|
collect('suspend_d','ts_code,trade_date,suspend_timing,suspend_type',trade_date='20260904')
|
||
|
|
collect('trade_cal','exchange,cal_date,is_open,pretrade_date',exchange='SSE',start_date='20260720',end_date='20260904')
|
||
|
|
for member in members:
|
||
|
|
code = member['con_code']
|
||
|
|
collect('daily','ts_code,trade_date,close,pre_close,pct_chg,vol,amount',trade_date='20260904',ts_code=code)
|
||
|
|
collect('moneyflow_dc','trade_date,ts_code,name,net_amount,net_amount_rate,pct_change,close',trade_date='20260904',ts_code=code)
|
||
|
|
collect('moneyflow','trade_date,ts_code,net_mf_amount',trade_date='20260904',ts_code=code)
|
||
|
|
result={'status':'collected','sector':'BK1147.DC','trade_date':str(target),'members':len(members),'snapshots':counts}
|
||
|
|
except Exception as exc:
|
||
|
|
result={'status':'partial','snapshots':counts,'error_type':type(exc).__name__,
|
||
|
|
'message':str(exc).replace(settings.tushare_token,'[redacted]')[:300]}
|
||
|
|
finally:
|
||
|
|
repository.close()
|
||
|
|
Path('../.trellis/tasks/09-06-capital-radar-daily-detail/research/collection-result.json').write_text(json.dumps(result,ensure_ascii=False,indent=2))
|
||
|
|
print(json.dumps({k:v for k,v in result.items() if k!='snapshots'},ensure_ascii=False))
|
||
|
|
print('Stored snapshots:',len(counts))
|
||
|
|
|
||
|
|
if __name__ == '__main__':
|
||
|
|
main()
|