金融行业里有个东西,叫"数据即服务"。我今年把这句话落地成了一个具体的项目,代号就叫 financial-services,一个面向个人投资复盘和家庭资产跟踪的数据聚合服务系统。不是那种花哨的炒股软件,而是完全自己掌控数据源、自己定义计算口径的后台服务。这几个月跑下来,从最初手动抓数据、拼表格,到现在每日自动生成投资日报、异常波动提醒、资产分布统计,整个流程算是彻底走通了。
这篇就当作一个阶段性的项目复现笔记,把整个设计思路、核心代码逻辑、踩过的坑,都记录下来。如果你也在琢磨怎么把手头的财务数据管起来,或者想给自己搞一套自动化的数据分析管道,这里面大部分思路是可以直接搬走用的。
1. 需求拆解与整体方案设计
1.1 原始需求到底要解决什么问题
刚开始我想做的很简单:每天收盘之后,不用自己一个个打开App去记账,能自动生成一份当天账户的净值变动、持仓盈亏和风险指标摘要。更进一步,我还想回答几个每年末都要倒腾半天的问题:今年总收益是多少、最大回撤发生在哪段、哪个板块贡献了主要收益。
这些事用Excel也能做,但问题是数据散落在各个地方。基金账户一个平台,股票账户另一个券商,银行理财又在一个App里。每条数据都有不同的导出格式,手工合并一次要花掉至少大半天,而且对不上账的情况经常发生。
所以 financial-services 的核心定位就是三件事:
- 统一数据采集:对接不同来源的行情和交易记录,做标准化清洗。
- 指标自动计算:市值、成本、收益、波动率、回撤这些指标,由系统统一按同一套算法跑。
- 定时输出结果:每日跑批生成报告,异常时主动推消息。
说白了,它就是一个私人的、自托管的金融数据中台。
1.2 技术选型背后的取舍逻辑
选型这件事,我卡了差不多两周。最后定下来的组合是 Python + PostgreSQL + Redis + Superset,跑在一台2核4G的小服务器上。为什么这么选,说一下当时的考量。
第一,语言用了 Python。这个基本没什么悬念。数据处理有 pandas、numpy,跑批任务有 APScheduler,行情接口有现成的 akshare 这类库可以兜底。就算有些数据源需要自己解析,Python 写起来也快。
第二,存储用了 PostgreSQL。金融数据有几个特点:字段结构相对固定、按时间序列查询多、对事务要求高。PostgreSQL 在这三个方向上都很均衡,而且内置了强大的日期时间函数和窗口函数,算累计收益率、移动平均这类指标,直接写 SQL 就能搞定,不用在内存里来回搬运。
第三,Redis 用来做缓存和任务锁。每日拉取行情不必重复请求,任务调度需要保证同一时间内只有一个实例在跑。这两件事用 Redis 都非常顺手。
第四,可视化没有自己造轮子,直接用了 Superset。说实话,自己画图表不是不能做,但那会消耗大量时间在无意义的前端调整上。Superset 拖拽一下就能生成资产分布饼图、净值走势折线图,省下来的时间足够我把计算逻辑打磨得更扎实。
这套组合还有一个隐性好处:全是开源软件,没有授权费用问题,数据也完全留在自己的机器上,不用担心第三方平台的政策变化导致接口不可用。
2. 核心模块拆解与数据模型设计
2.1 数据采集层:多源接入与标准化清洗
数据采集层是整个系统的基础,如果这层塌了,上面再怎么算都是空中楼阁。
我这边行情数据的来源主要是公开的行情接口,交易记录则靠各平台导出的 CSV 或者手工录入。原始数据格式五花八门,所以第一步是定义一个统一的标准格式。
最终我用了这样一张"交易流水表"来收口所有数据:
| 字段 | 类型 | 说明 |
|---|---|---|
| trade_id | varchar(64) | 全局唯一交易ID |
| account_id | varchar(32) | 账户标识,如 stock_a / fund_b |
| asset_code | varchar(16) | 资产代码,如 600519 或 110022 |
| asset_name | varchar(64) | 资产名称 |
| trade_type | varchar(8) | buy / sell / dividend |
| trade_date | date | 交易发生日期 |
| trade_price | numeric(12,4) | 成交价格 |
| trade_volume | numeric(14,2) | 成交数量(份额) |
| trade_fee | numeric(12,2) | 手续费 |
| raw_source | varchar(16) | 原始数据来源 |
不管原始数据长什么样,进了这张表就必须是这个结构。券商导出的成交明细列名可能叫"成交均价",基金平台可能叫"净值",清洗逻辑里统一做字段映射。这个环节我踩的坑比较多,后面专门写一节。
行情数据我单独建了一张 daily_quote 表,记录每个资产每日的开高低收、成交量、复权因子等。设计上特别注意了两个点:
- 用 (asset_code, trade_date) 做唯一索引,避免重复数据。
- 只保留后复权价格。做收益计算时后复权最省心,不然分红除权会把收益率曲线搞得七零八落。
2.2 指标计算层:计算口径的"唯一事实来源"
指标计算层是整个系统最核心的部分。同一笔交易,不同人算出来的收益率可能不一样,原因就在于分母的分歧。所以我在设计时定了一套规则,写死在一个配置表里。
比如"持仓成本"的定义,系统里统一采用移动加权平均法:每次买入后,新成本等于(旧成本 * 旧数量 + 买入金额 + 手续费)除以(旧数量 + 买入数量)。卖出时不改变成本单价,只减少数量。分红则按实际情况区分为现金分红和再投资。
再比如"收益率",我同时维护了三种口径:
- 累计收益率:当前资产净值相对于总投入资金的变化幅度。
- 年化收益率:按实际持有天数做指数折算,公式是 (1 + 累计收益率) 的(365 / 持有天数)次方减一。
- 时间加权收益率:这个方法剔除了出入金的影响,每一段资金进出之间的区间收益单独计算,再连乘。用来评估"投资能力"时这个指标最干净。
这些口径定义好之后,所有下游报表都只认这套计算引擎的输出,不会出现同一个指标在不同图表里数字对不上的尴尬情况。
2.3 调度与通知:让系统主动说话
数据有了、指标也能算了,但如果不自动跑,那跟手动Excel也没什么区别。我引入了 APScheduler 做了三个定时任务:
- 每日 15:30 拉取当日行情数据,做增量更新。
- 每日 16:00 触发现金流和持仓的重算,生成当日的资产快照。
- 每日 16:30 汇总生成投资收益日报,并推送提醒到手机。
推送这块用的是 Server 酱的简单接口,本质上就是发一个 HTTPS 请求。日报内容包含当日总资产、日涨跌额、涨跌比例、持仓变动列表、风险提示等。系统的价值从这里开始显现:它能"主动说话",而不再是我去翻数据。
3. 实操过程与核心环节实现
3.1 数据库表结构与初始化脚本
先看实际的建表 SQL,整个过程我按模块逐步推进。第一个落地的是资产净值快照表 account_daily_value,这张表是后续所有报表的数据源头。
CREATE TABLE account_daily_value ( account_id varchar(32) NOT NULL, value_date date NOT NULL, total_asset numeric(14,2) NOT NULL, -- 总资产 available_cash numeric(14,2) NOT NULL DEFAULT 0, -- 可用现金 market_value numeric(14,2) NOT NULL DEFAULT 0, -- 持仓市值 accrued_income numeric(14,2) NOT NULL DEFAULT 0, -- 累计收益 PRIMARY KEY (account_id, value_date) );这张表的特殊之处在于,它存的是状态值,不是流水。每天的收盘后,系统根据当日所有持仓的收盘价重新计算一次市值,加上现金余额,得到总资产。这样看历史某一天的资产规模,直接按日期查就好了,不需要从流水重算,性能快得多。
为了确保幂等,写入时用 UPSERT 逻辑:如果某一天的记录已存在,就更新而不是插入新行。这样即使调度任务重复执行,结果也是一致的。
3.2 行情拉取与增量更新逻辑
行情数据拉取我喜欢用"日期增量"策略。一个资产从上市到现在的全部历史行情,只有在首次引入时才全量拉一遍,之后就只补最近几天的数据。
import akshare as ak import pandas as pd from sqlalchemy.dialects.postgresql import insert def update_quote(code: str, start_date: str): # 拉取指定区间的前复权行情 df = ak.stock_zh_a_hist(symbol=code, start_date=start_date, adjust="qfq") records = [] for _, row in df.iterrows(): records.append({ "asset_code": code, "trade_date": row["日期"], "open_price": row["开盘"], "close_price": row["收盘"], "high_price": row["最高"], "low_price": row["最低"], "volume": row["成交量"], }) stmt = insert(daily_quote).values(records) stmt = stmt.on_conflict_do_update( constraint="uq_asset_quote_date", set_={"close_price": stmt.excluded.close_price, "volume": stmt.excluded.volume} ) db.execute(stmt) db.commit()这里有个细节需要注意:start_date不能只传当天的日期。比如今天补数时,可能前两天因为数据源异常漏拉了,导致中间有空档。我后来干脆改成:每次拉取前,先查一下库里这个代码的最大日期,然后从最大日期往前倒推3天开始拉。多出来的重复数据由唯一索引去重,这样既保证了数据连续性,代码也很简单。
建议先在数据库里创建唯一索引:
CREATE UNIQUE INDEX uq_asset_quote_date ON daily_quote(asset_code, trade_date);3.3 持仓成本与收益计算
接下来是关键的持仓成本计算。这个模块我重写过三遍,前两版在遇到分批买入、部分卖出、分红再投资时都会对不上账。最后改成用 SQL 窗口函数逐笔滚动计算,才稳定下来。
核心思路是:先把一个账户的所有交易流水按时间排序,然后逐行计算"本次变动后的累计持仓数量"、"本次变动后的持仓成本单价"和"本次交易产生的盈亏"。
WITH ordered_trades AS ( SELECT asset_code, trade_date, trade_type, trade_price, trade_volume, trade_fee, CASE WHEN trade_type = 'buy' THEN trade_volume WHEN trade_type = 'sell' THEN -trade_volume ELSE 0 END as signed_volume, CASE WHEN trade_type = 'buy' THEN trade_price * trade_volume + trade_fee WHEN trade_type = 'sell' THEN 0 ELSE 0 END as buy_amount FROM trade_records WHERE account_id = :acct AND trade_date <= :end_date ORDER BY trade_date, trade_id ) SELECT asset_code, SUM(signed_volume) as hold_volume, CASE WHEN SUM(CASE WHEN signed_volume > 0 THEN signed_volume ELSE 0 END) = 0 THEN 0 ELSE SUM(buy_amount) / NULLIF(SUM(CASE WHEN signed_volume > 0 THEN signed_volume ELSE 0 END), 0) END as avg_cost FROM ordered_trades GROUP BY asset_code;这个 SQL 做了简化处理,实际在计算"移动加权成本"时,卖出不会更新成本单价,所以总买入金额除以历史累计买入数量,得到的就是当前持仓的平均成本。这个方法有一个前提:没有把卖出盈亏再摊回成本里,这正是符合直觉的"持仓成本"含义。
实现之后,我又做了一层校验逻辑:用系统算出的持仓市值,去跟券商App上的市值对账,误差超过0.5%时自动标记异常。这个校验功能救了很多次场,后面会详细说。
3.4 最大回撤与波动率计算:风险指标的落地
收益指标好算,风险指标才是容易出错的地方。最大回撤的定义是"任一高点买入后到后续最低点的最大亏损幅度",用 SQL 写窗口函数比较方便。
计算逻辑分两步:
- 计算每天的累计净值(先把每日总资产序列做归一化)。
- 用窗口函数计算截至当日的"历史最高净值",回撤率等于(当前净值 / 历史最高净值 - 1)。
WITH daily_values AS ( SELECT value_date, total_asset, MAX(total_asset) OVER (ORDER BY value_date) as peak_asset FROM account_daily_value WHERE account_id = :acct ) SELECT value_date, total_asset, peak_asset, (total_asset / peak_asset - 1) as drawdown_rate FROM daily_values ORDER BY drawdown_rate ASC LIMIT 1;注意这里有个细节:峰值是"截至当时"的峰值,不是整个区间的峰值。如果直接用全局最大值,那么回撤曲线的起点会被错误地归零,计算出的回撤区间不准确。
波动率的计算我用的是近20个交易日收益率的标准差,再年化:年化波动率 = 日波动率 * sqrt(250)。这个数在日报里作为风险参考指标输出,数值本身不复杂,但计算时要用前面提过的时间加权收益率而不是简单收益率,否则遇到出入金时指标会失真。
3.5 仪表盘与自动化报表
最后一步是把算好的数据可视化。Superset 连接 PostgreSQL 之后,我建了几个核心图表:
- 总资产走势图(按账户分组,时间序列折线图)。
- 资产分布饼图(按资产大类聚合当前市值)。
- 月度盈亏柱状图(按月聚合已实现+未实现收益)。
- 回撤区间图(高亮显示历史上最大回撤发生的区间)。
Superset 只需要在 SQL Lab 里写查询语句,拉到图表界面配置即可。我最常用的查询是按月份聚合收益:
SELECT DATE_TRUNC('month', value_date) AS month, SUM(accrued_income) AS income FROM account_daily_value WHERE account_id = 'main' GROUP BY 1 ORDER BY 1;报表推送是最后一块拼图。每日收盘后,系统把当天核心指标拼成一段纯文本,推送到手机。实际推送效果大致是:
【投资日报】2024-11-20 总资产:228,503.18(+0.82%) 持仓市值:184,220.50 现金:44,282.68 今日盈亏:+1,856.32 本月累计:+4,231.08 风险提示:暂无这个日报的价值在于是"被动获取"的,不用主动去查,反而更容易坚持记录。人都是有惰性的,系统自动化最大的收获就是帮助用户养成了每日审视资产的习惯。
4. 实际运行中遇到的问题与排查实录
4.1 数据源接口字段频繁变动
第一个大坑是第三方数据库的字段名不稳定。最开始写好的采集代码,运行三周后忽然报错,原因是源接口把"收盘"字段从中文改成了英文,或者干脆调整了列顺序。这类错误隐蔽性很强,因为接口调用本身不报错,只是返回的 DataFrame 结构变了。
解决方案分两层:
- 采集后立即做结构校验,检查必须包含的字段是否齐全,不齐全就报警,而不是直接进入清洗流程。这样错误能在最早的时间点暴露。
- 字段映射走字典配置,不要硬编码列名。哪怕源字段变了,改一行配置就能恢复。
4.2 分红除权导致收益率跳变
这个问题非常经典。持仓的股票分红了,如果直接按照除权后的行情价格计算市值,会发现资产"凭空少了一笔",但账户余额里又多了现金分红,如果只盯着市值曲线看就会觉得是亏损。
处理方法是统一使用后复权价格计算收益。后复权价格会把历史价格按分红、送股做调整,这样算出的收益率曲线不会因为除权除息而产生虚假的下跌。
但这里我又踩了一个更深的坑:基金的分红再投资和股票的后复权逻辑并不完全一样。基金如果选了现金分红,净值不变,但持有份额不变、现金增加,收益按增加现金算。如果选了红利再投,持有份额增加,成本也相应增加。所以在 trade_records 表中,我用 trade_type = 'dividend' 来标识分红事件,并结合 dividend_action 字段区分两种处理方式。这块分类逻辑不处理好,总资产的连续性校验一定过不了。
4.3 时区问题导致的日期错位
调度任务使用了服务器本地时间,而行情数据源返回的时间往往是东八区时间。如果服务器时区没设置一致,就会出现"今天的数据被记到昨天"的问题。
排查过程不算复杂,但足够折腾。现象是当天日报显示的总资产比券商App少了约2%,而且对账时发现少了当日持仓市值。
最终解决方式是在所有涉及日期计算的位置,统一强制时区:
from datetime import datetime, timezone, timedelta CN_TZ = timezone(timedelta(hours=8)) today = datetime.now(CN_TZ).date()调度的 cron 表达式也改成按东八区设定,而不是依赖系统的默认时区。
4.4 报表与券商数据对不平
最后谈一下对账问题。系统算出来的市值,跟券商App显示的不一致,这是我在调试期间最常遇到的问题。出错点往往不在计算本身,而在数据采集的准确性上。
尤其要注意的是,券商App显示的"当日盈亏"和"持仓盈亏"口径不同:当日盈亏用的是当日收盘价对比昨日收盘价,持仓盈亏用的是当前价对比买入成本价。如果系统算出的结果和App对不上,先确认是不是在拿苹果跟橘子比。
我加了一个对账任务:每天收盘后,自动把系统内持仓的市值按代码汇总,并跟手动导入的券商持仓快照做比对。差额超过0.5%的资产会被标记出来,第二天我去查这条资产的成交记录。几次排查后发现,问题集中在分红再投资的份额更新不及时、或者手续费漏记了。对账机制虽然简单,但能快速锁定错误范围,大大缩短排查时间。
5. 一些值得记录的实测心得
跑了大半年,整体系统稳定之后,我有几条比较深刻的体会想单独写出来。
第一,自建数据服务的核心价值不在"算",而在"攒"。行情数据、交易数据、账户净值数据,这些数据积累得越久,能算出来的指标越有价值。三个月的数据只能看个寂寞,三年的数据才能看出真实的波动特征和风险偏好。所以系统设计时,把"数据连续性和完整性"放在了最高优先级,宁可少一个指标,也不能丢一天的数据。
第二,指标口径的统一比算法本身更重要。很多人做数据分析卡在"算法不够高级",但实际生产环境中,更多事故发生在"昨天算的口径跟今天不一样"。我在配置中心用一份 YAML 文件维护了所有指标的计算口径说明和参数定义,任何计算引擎的改动都必须同时更新这份文档。这看起来不起眼,但在半年后回顾的时候,作用非常大。
第三,自动化并不是"越多的自动越好"。我一开始想做的很全,佣金试算、止损提醒、自动调仓建议,什么都往系统里塞。结果发现,自动止损提醒这类功能在情绪上很容易被忽略,而且误报率一高,人对系统的信任度会迅速下降。最后我把自动交易相关的功能全部砍掉,只保留了信息聚合、报表、风险指标汇总,把决策完全留给人。系统的角色回归到"数据服务"本身,反而更顺了。
如果你想在自己的服务器上复刻这套系统,我建议不要一开始就追求大而全。从最简单的每日净值快照开始,跑两周稳定了,再加行情自动更新;再跑两周,再加回撤计算;一层层往上叠。数据管道这种东西,迭代过程中要动数据结构,早期的简单设计反而让重构成本更低。
以上就是 financial-services 这个项目从零到一的核心过程。目前这套服务还在持续跑着,每天定时出报告,数据已经在慢慢积累。回头再看,最难的不是编码,而是把金融业务概念准确地翻译成工程上的数据模型和计算规则。这个过程很磨人,但磨完之后,收益也是长期的。