7f93d6b0f5
- Introduced ActiveMoneyflowSource to fetch optional active-order flow, enhancing the sector radar's data capabilities. - Updated StockFactRecord and DailyAggregateRecord to include pct_change and active_buy_net_amount_yuan for improved financial insights. - Modified the build process to incorporate active moneyflow data without invalidating main rankings on failure. - Enhanced the HTTP API to return detailed sector history and metrics, including pct_change and active buy metrics for members. - Updated tests to validate the new functionality and ensure data integrity across various scenarios.
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()
|