从接口获取到数据资产:本地化股票数据的工程化落地实践
做了好几年股票数据相关的开发,从最开始在浏览器里手动查行情,到调用零散的API接口,再到今天有一套完整的本地化数据资产体系,中间踩的坑、改的架构、写的代码,说起来都是故事。这篇文章我想做一个系统的总结:如何把零散的API调用整合成可维护、可扩展、可复用的数据资产。我会从架构设计、分层实现、工程规范三个维度展开,分享一套适合个人开发者和小团队的股票数据工程化落地实践。
问题的起源:零散接口的痛点
最开始我的股票数据获取方式非常原始:写一个Python脚本,需要什么数据就调什么接口,行情用实时接口、K线用历史接口、资金用资金接口,每个脚本都是独立的,没有任何公共逻辑。这种方式在项目初期跑得很欢,但随着需求增多,问题越来越严重。
第一个痛点是重复造轮子。每个脚本都要自己处理接口鉴权、请求重试、数据解析、异常处理。改一个解析逻辑要改所有脚本,维护成本随脚本数量线性增长。
第二个痛点是数据不一致。不同脚本独立拉取数据,同一只股票的同一个指标可能来自不同时间的接口调用,导致数据对不上。比如一个脚本上午拉了某股票的资金流向,另一个脚本下午拉了同一只股票的实时行情,两者之间的关联分析就会出问题。
第三个痛点是无法复用。我写的K线拉取逻辑是"一次性"的,换一个项目就得重写。没有标准化的数据格式和接口,数据资产无法积累。
第四个痛点是缺乏监控。接口调用失败、数据异常、存储出错,这些问题只能靠用户反馈或者手动检查,没有自动化的监控和告警。
这些问题逼着我思考:能不能把零散的接口调用整合成一个统一的数据框架?让数据从采集到存储到服务都有规范,让每一份数据都成为可追溯、可复用的资产?
架构设计:四层数据架构
经过几轮迭代,我最终确定了四层架构:数据采集层、数据清洗层、数据存储层、数据服务层。每一层职责清晰,通过标准格式的JSON数据传递,层与层之间解耦。
┌─────────────────────────────────────────────┐ │ 数据服务层 (Service) │ │ 对外提供统一查询API、数据分析接口 │ ├─────────────────────────────────────────────┤ │ 数据存储层 (Storage) │ │ SQLite/DuckDB/JSON多引擎、版本管理 │ ├─────────────────────────────────────────────┤ │ 数据清洗层 (Clean) │ │ 格式标准化、异常值处理、数据校验、去重 │ ├─────────────────────────────────────────────┤ │ 数据采集层 (Collector) │ │ 接口路由、并发控制、重试机制、限流 │ └─────────────────────────────────────────────┘数据采集层:统一的接口路由
采集层的核心设计是一个统一的接口路由器。所有数据请求都走这个路由器,由它决定调用哪个接口、用什么参数、并发度多少、失败后如何重试。
首先定义一个标准的请求协议:
fromdataclassesimportdataclassfromtypingimportList,Optional@dataclassclassDataRequest:data_type:strdm_list:List[str]params:dictpriority:int=0timeout:int=30@dataclassclassDataResponse:data_type:strdm:strraw_data:dicttimestamp:strsuccess:boolerror_msg:str=""然后是接口路由器,根据data_type分发到不同的接口:
classInterfaceRouter:INTERFACE_MAP={"stock_list":{"path":"/base/gplist","method":"GET"},"realtime":{"path":"/time/real/{dm}","method":"GET"},"kline":{"path":"/time/history/trade/{dm}/{level}","method":"GET"},"finance":{"path":"/time/f10/fi/{dm}","method":"GET"},"l2_sign":{"path":"/time/real/trace/l2sign/{dm}","method":"GET"},"onebyone":{"path":"/time/real/trace/onebyone/{dm}","method":"GET"},"zlzjzs":{"path":"/time/zijin/zlzjzs/{dm}","method":"GET"},"zjlrqs":{"path":"/time/zijin/zjlrqs/{dm}","method":"GET"},"longhubang":{"path":"/time/data/longhubang","method":"GET"},"bshgt":{"path":"/time/data/bshgt","method":"GET"},}def__init__(self,rate_limiter=None):self.rate_limiter=rate_limiterorRateLimiter(5)self.session=requests.Session()self.session.headers.update({"User-Agent":"Mozilla/5.0","Referer":"https://quote.example.com"})defroute(self,request:DataRequest)->List[DataResponse]:config=self.INTERFACE_MAP.get(request.data_type)ifnotconfig:returnself._error_response(request,"Unknown data_type")results=[]fordminrequest.dm_list:self.rate_limiter.acquire()url=self._build_url(config["path"],dm,request.params)try:resp=self.session.get(url,timeout=request.timeout)resp.raise_for_status()raw=resp.json()results.append(DataResponse(data_type=request.data_type,dm=dm,raw_data=raw,timestamp=datetime.now().isoformat(),success=raw.get("rc")==0))exceptExceptionase:results.append(DataResponse(data_type=request.data_type,dm=dm,raw_data={},timestamp=datetime.now().isoformat(),success=False,error_msg=str(e)))returnresults这个路由器的好处是:所有接口调用逻辑集中在一处,加新接口只需要在INTERFACE_MAP里加一行;统一的异常处理和重试机制;并发控制和限流策略全局生效。
数据清洗层:标准化与校验
采集层拉回来的原始数据是五花八门的,有的字段名不一样,有的数据类型不一致,有的有脏数据。清洗层的职责就是把这些数据标准化成统一的格式。
清洗层的核心是一个字段映射表和校验规则集:
FIELD_MAPPING={"realtime":{"f43":"cjjg","f47":"cjl","f58":"cjsj","f169":"jyzd",},"zlzjzs":{"f62":"zlJlr","f184":"zlJlb","f70":"shJlb",},"zjlrqs":{"f78":"f5MinZlJe",},"l2_sign":{"f164":"ddx","f165":"ddy","f166":"ddz","f167":"ddf",},}VALIDATION_RULES={"cjjg":lambdax:x>0,"cjl":lambdax:x>=0,"dm":lambdax:len(x)==6andx.isdigit(),"cjsj":lambdax:len(x)>=8,}classDataCleaner:def__init__(self):self.field_mapping=FIELD_MAPPING self.validation_rules=VALIDATION_RULESdefclean(self,response:DataResponse)->Optional[dict]:ifnotresponse.success:returnNonedata=response.raw_data.get("data",{})mapping=self.field_mapping.get(response.data_type,{})cleaned={"dm":response.dm,"data_type":response.data_type}forsrc_field,dst_fieldinmapping.items():value=data.get(src_field)ifvalueisnotNone:cleaned[dst_field]=self._coerce(dst_field,value)else:cleaned[dst_field]=Noneforfield,ruleinself.validation_rules.items():iffieldincleanedandcleaned[field]isnotNone:ifnotrule(cleaned[field]):cleaned["_invalid"]=Truebreakreturncleaneddef_coerce(self,field,value):try:iffieldin("cjjg","zlJlr","f5MinZlJe"):returnfloat(value)eliffieldin("cjl",):returnfloat(value)else:returnstr(value).strip()except(ValueError,TypeError):returnNone清洗层的设计原则是:每个字段都要经过映射表转换,确保输出格式统一;每个关键字段都有校验规则,不合格的数据标记为无效但不丢弃;清洗过程是无状态的,便于并行。
数据存储层:多引擎与版本管理
存储层支持多种存储引擎,不同类型的数据可以存在不同引擎里。我当前的配置是:基础数据用SQLite,K线数据用DuckDB,配置类数据用JSON文件。
classDataStorage:def__init__(self,config):self.engines={}forengine_name,engine_configinconfig.items():self.engines[engine_name]=self._create_engine(engine_config)def_create_engine(self,config):ifconfig["type"]=="sqlite":returnSQLiteEngine(config["path"])elifconfig["type"]=="duckdb":returnDuckDBEngine(config["path"])elifconfig["type"]=="json":returnJSONEngine(config["path"])raiseValueError(f"Unknown engine type:{config['type']}")defwrite(self,data_type:str,records:List[dict]):engine_name=self._resolve_engine(data_type)engine=self.engines[engine_name]engine.write(records)defread(self,data_type:str,query:dict)->List[dict]:engine_name=self._resolve_engine(data_type)engine=self.engines[engine_name]returnengine.read(query)def_resolve_engine(self,data_type:str)->str:mapping={"stock_list":"sqlite","realtime":"sqlite","kline":"duckdb","finance":"duckdb","zlzjzs":"duckdb","l2_sign":"duckdb","config":"json",}returnmapping.get(data_type,"sqlite")版本管理在存储层实现。每次写入都带一个版本号,版本号由时间戳生成。查询时默认返回最新版本,但可以指定历史版本。
classVersionedStorage(DataStorage):defwrite_with_version(self,data_type,records,version=None):ifversionisNone:version=datetime.now().strftime("%Y%m%d_%H%M%S")forrinrecords:r["_version"]=version self.write(data_type,records)self._record_version(data_type,version,len(records))defread_by_version(self,data_type,query,version=None):ifversion:query["_version"]=versionreturnself.read(data_type,query)数据服务层:对外的统一API
服务层是数据资产的门面,对外提供统一的查询接口。内部可以是Web API、命令行工具、或者Python包。我实现的是一个Python包,直接import使用,也可以通过FastAPI包装成Web服务。
classDataService:def__init__(self,storage:VersionedStorage):self.storage=storagedefget_stock_list(self,filter_by=None)->List[dict]:data=self.storage.read("stock_list",{})iffilter_by:forkey,valueinfilter_by.items():data=[dfordindataifd.get(key)==value]returndatadefget_kline(self,dm:str,start:str,end:str,level:str="101")->List[dict]:returnself.storage.read("kline",{"dm":dm,"cjsj_start":start,"cjsj_end":end,"level":level,})defget_realtime(self,dm:str)->dict:results=self.storage.read("realtime",{"dm":dm})returnresults[0]ifresultselseNonedefget_zijin(self,dm:str,data_type:str="zlzjzs")->dict:results=self.storage.read(data_type,{"dm":dm})returnresults[0]ifresultselseNonedefbatch_query(self,queries:List[dict])->List[dict]:results=[]forqinqueries:data=self.storage.read(q["data_type"],q.get("filter",{}))results.extend(data)returnresults工程规范:让数据资产可维护
架构搭好只是第一步,真正让数据资产可持续运营的是工程规范。我在实践中总结了几条核心规范:
每个数据类型都有唯一的data_type标识,贯穿采集、清洗、存储、服务全流程。新增数据类型时,先定义data_type,然后逐层实现。
每个字段都有标准名称和数据类型。dm永远是字符串类型的6位股票代码,cjjg永远是浮点类型的成交价格,cjsj永远是字符串类型的时间戳。命名规范在FIELD_MAPPING里集中定义。
每次数据操作都有日志记录。包括拉取了多少只股票、写入了多少条记录、失败了多少条、耗时多久。这些日志用于监控和审计。
定期做数据质量巡检。每周跑一次全量校验,检查数据完整性(有没有缺失的日期)、一致性(不同数据源的同一指标是否一致)、准确性(价格是否合理)。
配置与代码分离。所有接口地址、存储路径、并发参数等都放在配置文件里,代码只读取配置,不硬编码。
CONFIG={"api":{"base_url":"https://api.example.com","rate_limit":5,"timeout":15,},"storage":{"sqlite":{"type":"sqlite","path":"stock.db"},"duckdb":{"type":"duckdb","path":"stock.duckdb"},"json":{"type":"json","path":"config/"},},"schedule":{"daily_update":"15:30","quality_check":"02:00",},}落地效果:从混乱到有序的转变
做完这套架构之后,最直观的变化是新项目的启动时间。之前做一个新的分析需求,从获取数据到跑通流程可能要一周,现在只需要一天——直接用DataService,调几个接口就能拿到数据。
第二个变化是数据质量。因为有了统一的清洗和校验,数据的一致性和准确性大幅提升。之前做跨指标分析经常出现数据对不上的情况,现在很少遇到了。
第三个变化是可维护性。之前修改一个字段的解析逻辑,要改五六个脚本;现在只需要改FIELD_MAPPING里的一行。代码量减少了,bug也少了。
第四个变化是可扩展性。想加新的数据源?在InterfaceRouter的INTERFACE_MAP里加一行,在DataCleaner的FIELD_MAPPING里加一行,就能接入。想加新的存储引擎?在DataStorage里加一个Engine实现,就能切换。
总结与展望
从零散的接口调用到成体系的数据资产,这条路走了差不多两年。中间最大的体会是:数据工程的价值不在技术选型,而在架构设计和工程规范。技术选型可以迭代,但架构一旦定下来,后续的演进成本就取决于当初的设计质量。
如果你正在从零开始搭建股票数据系统,我的建议是:先把四层架构搭起来,哪怕每一层的实现都很简陋;然后在使用过程中逐步补全细节,而不是一开始就追求完美。股票数据的特点是变化慢,架构的使用寿命很长,前期投入的时间会在后续几年里持续产生回报。
未来我计划在这套架构上增加更多智能能力:用Agent做自动数据巡检和异常检测,增加数据血缘追踪(每条数据的来源和处理过程可追溯),以及基于数据质量的自动评分。这些都是在现有架构上可以平滑扩展的方向。
接口说明
| 接口路径 | 用途 | 核心参数 | 核心返回字段 |
|---|---|---|---|
| base/gplist | 获取全市场股票列表 | - | dm, mc, hy, ssrq |
| time/real/{dm} | 获取实时行情 | dm=股票代码 | f43(现价), f47(成交量), f58(时间), f169(方向) |
| time/history/trade/{dm}/{level} | 获取历史K线 | dm=代码, level=周期 | klines(时间,开,收,高,低,量) |
| time/f10/fi/{dm} | 获取财务指标 | dm=股票代码 | 营收, 净利润, ROE等 |
| time/real/trace/l2sign/{dm} | 获取L2指标 | dm=股票代码 | ddx, ddy, ddz, ddf |
| time/real/trace/onebyone/{dm} | 获取逐笔交易 | dm=股票代码 | cjsj, cjjg, cjl, jyzd |
| time/zijin/zlzjzs/{dm} | 获取资金走势 | dm=股票代码 | zlJlr, zlJlb, shJlb |
| time/zijin/zjlrqs/{dm} | 获取资金趋势 | dm=股票代码 | f5MinZlJe等 |
| time/data/longhubang | 获取龙虎榜数据 | 日期 | 营业部, 买入金额, 卖出金额 |
| time/data/bshgt | 获取北向资金数据 | 日期 | 沪股通, 深股通净流入 |
资料参考:ig50.com