Pathway 流式计算期权 Greeks:Option Greeks 示例项目架构、源码与运行全解析
本示例以 Pathway 实时数据框架(Live Data Framework) 为核心,从 Databento 市场数据服务拉取期权合约定义与逐档行情,在流式数据上实时求解隐含波动率并计算 Delta、Gamma、Theta、Vega、Rho 五大期权 Greeks,最终通过 Streamlit 仪表盘与表格查询服务进行展示。读完本文你将掌握该示例的完整运行流程、每一段 Pipeline 的底层实现原理,以及如何将同一套计算链复用到你自己的实时风控或量化分析场景。
项目要解决的问题:为什么需要实时计算 Option Greeks
期权(Options)是一种赋予持有者在约定期间内以指定价格买卖标的资产的权利(而非义务)的金融衍生品。Option Greeks(期权希腊字母)是评估期权价格对各种风险因素敏感度的量化指标,交易员与风险管理经理正是借助它们做出更明智、更具策略性的决策,从而增强风险管理能力并优化收益。
在本示例中,Greeks 的计算完全由 Pathway 的数据流框架完成,行情数据由 Databento API 提供,计算结果则呈现在一个 Streamlit 仪表盘中。整套方案最有价值的点在于:计算不是对一份静态 CSV 跑一次离线脚本,而是随着行情持续到达不断增量更新——这正是 Pathway 作为流处理与实时分析框架的典型应用场景。
整体工作流如下图所示(来源:option-greeks.svg):
示例项目全景与文件组织
本项目源码位于仓库的 examples/projects/option-greeks/ 目录下,共提供两种数据来源的等价实现:
examples/projects/option-greeks/
├── data/
│ ├── definition.csv # 静态模式使用的期权合约定义数据
│ └── options.csv # 静态模式使用的订单簿(mbp-1)行情数据
├── .env # 存放 Databento API Key 等运行参数
├── greeks-replay.py # 通过 Databento Historical API 回放历史行情并实时计算 Greeks
├── greeks-static.py # 读取本地 CSV 静态数据的等价脚本
├── querying.py # FastAPI 表格查询服务:把 Pathway 输出表暴露为 HTTP 接口
├── requirements.txt # Python 依赖清单
├── streamlit_ux.py # Streamlit 用户界面(仪表盘)
├── option-greeks.svg # 架构/工作流示意图
└── Streamlit.png # 运行结果截图
各文件职责如下:
data/:greeks-static.py读取的样例数据。definition.csv 含约 3250 条期权合约定义,options.csv 含约 1.8 万条逐笔档位行情;.env:配置 Databento API Key 与 HTTP 服务端口等参数;greeks-replay.py:调用 Databento Historical API 拉取历史行情,模拟"实时回放";Greeks 在数据流上持续计算;greeks-static.py:与 replay 等价的计算逻辑,但数据来自本地 CSV(便于离线调试);querying.py:将 Pathway 输出表注册为可查询别名,并启动 FastAPI 服务,提供/get_table?alias=...接口;streamlit_ux.py:Streamlit 仪表盘,轮询上面的 HTTP 接口并把 Greeks 渲染成表格;requirements.txt:运行两个脚本所需的全部 Python 库。
值得注意:上面的依赖/交互关系中,querying.py 是 streamlit_ux.py 与两个计算脚本共同依赖的"表 → Web 接口"桥接模块,README 目录树未列出它,但它实际参与了整套系统的运行。
准备运行环境
安装依赖库
在项目目录(examples/projects/option-greeks/)下执行:
pip install -r requirements.txt
requirements.txt 中的依赖与其用途如下:
| 依赖库 | 用途 |
|---|---|
pathway |
流式数据计算引擎(表、join、udf、订阅等核心能力) |
databento |
访问 Databento 市场数据 API |
pandas |
中间数据整理与 Streamlit 侧 DataFrame 渲染 |
scipy |
正态分布函数(norm.cdf/pdf)与隐含波动率求根(root_scalar) |
python-dotenv |
读取 .env 中的配置 |
streamlit |
仪表盘 UI |
fastapi / uvicorn / pydantic |
表格查询服务与 ASGI 运行器 |
配置 Databento API Key
首先需要在 Databento 平台注册账号并创建 API Key(新账户通常可获得免费额度)。随后编辑 .env 填入你的密钥和其他参数:
API_KEY =<在此填写你的 Databento API Key>
PORT =<在此填写端口号,README 示例中为 16001>
其中 API_KEY 被两个计算脚本通过 os.environ.get("API_KEY") 读取;PORT 则被脚本与 UI 共同用于暴露 Pathway 表格查询服务,README 中的验证地址 http://localhost:16001/get_table?alias=table_greeks 正是基于该端口。
认识样例数据:合约定义与逐笔行情
两套脚本都依赖两类输入,理解其字段含义是读懂 Pipeline 的前提(字段注释直接来自两个计算脚本中的 Schema 声明,见 greeks-static.py)。
合约定义数据(静态 CSV 头部示例):
ts_recv,raw_symbol,expiration,instrument_class,strike_price,underlying,instrument_id,time,diff
1712275200000000000,"ESZ4 C5100",1734705000000000000,"C",5100.0,"ESZ4",434835,1721398215586,1
定义表中每条记录对应一个期权合约,关键字段包括:
ts_recv:数据接收时间(纳秒级时间戳);raw_symbol:期权的原始交易符号;expiration:期权到期时间(纳秒级时间戳),用于换算距离到期的时间;instrument_class:期权类型,"C"代表 Call(看涨期权)、"P"代表 Put(看跌期权);strike_price:行权价(示例数据中原始值放大了 1e9 倍,脚本会做归一化);underlying:标的合约符号(如ESM4、ESZ4);instrument_id:期权合约标识符,Dashboard 中用它标识每一行。
行情数据(mbp-1 档位快照,CSV 头部示例):
ts_recv,symbol,bid_px_00,ask_px_00,raw_symbol,bid_px,ask_px
2024-04-04 22:00:00.012859335+00:00,ESM4,5204.25,5190.0,ESM4,5204.25,5190.0
行情表主要关心三列:raw_symbol(标的/期权符号)、bid_px(买价)、ask_px(卖价)。代码只把符合 Schema 声明的列读入表格,CSV 中多出的辅助列不会影响计算。需要注意的是行情数据原始单位为 1e-9,脚本读入后统一除以 1e9 转换为常规价格单位。
主程序一:greeks-replay.py——历史行情回放驱动的实时 Greeks
greeks-replay.py 是完整版实现,它从 Databento Historical API 拉取"历史快照"并按时间节奏回放,从而在 Pathway 中模拟实时行情。
输入 Schema 与自定义 Connector
脚本先用 pw.Schema 声明定义表结构,再通过 pw.io.python.ConnectorSubject 自定义数据源:
class DefinitionInputSchema(pw.Schema):
ts_recv: int # 数据接收时间(ns)
raw_symbol: str # 期权符号
expiration: int # 到期时间
instrument_class: str # 期权类型(C/P)
strike_price: float # 行权价
underlying: str # 标的合约
instrument_id: int # 期权标识
ConnectorSubject 的 run() 中调用 client.timeseries.get_range(...) 从 Databento 拉取 definition 架构数据,逐行提取需要的属性并通过 self.next(...) 发射到 Pathway 表格(greeks-replay.py)。相关 Databento 参数:
dataset = "GLBX.MDP3":CME Globex MDP 3.0 数据集;schema = "definition":拉取合约定义;symbols = ["ES.OPT"]:通过stype_in=db.SType.PARENT解析出所有以 ES 为根符号的期权;start/end:设定历史区间(示例为 2024-04-04 开始的一天);- 定义数据本身变化频率低,因此循环内不做
time.sleep(注释明确说明"as this data is static")。
行情(价格)输入侧同样用 pw.io.python.read 读取 mbp-1(最优买卖档)数据。关键细节包括:借助 Databento 的 InstrumentMap 把 instrument_id 解析回符号;当买/卖价等于 (1<<63) 附近的标记值时视为"未知价格"并跳过;价格统一除以 1e9 归一化(greeks-replay.py)。
数据流上的计算管线
完成数据接入后,全部计算都用 Pathway 表操作表达:
- 过滤合约定义:先
filter(pw.this.underlying == front_month_symbol)只保留近月合约ESM4的定义,再过滤出instrument_class为"C"或"P"的行; - 提取需要询价的符号集合:通过
reduce(symbol_tuple=pw.reducers.tuple(...))聚合出所有期权符号,再经pw.debug.table_to_pandas(实现见 python/pathway/debug/init.py)拿到符号列表,作为第二段行情请求的symbols参数; - 计算期权中间价:对行情表按
raw_symbol做groupby().reduce(),用pw.reducers.avg对买价/卖价求均值,得到option_midprice = (avg(bid) + avg(ask)) / 2; - Join 两张表:将"合约定义"与"期权中间价"按
raw_symbol等值连接(join),产出每个期权最新的中间价; - 取期货价格作为远期价格 F:调用
table_mbp1.ix_ref(front_month_symbol).option_midprice从行情表中按主键取出ESM4期货自身的中间价,作为 Black 模型中的远期价格F。ix_ref是 Pathway 提供的"按键引用某一行字段"的表操作原语(实现见 python/pathway/internals/table.py); - 计算距离到期时间(年):用
@pw.udf装饰的纯 Python 函数把纳秒时间戳换算成以年为单位的时间T,即(expiration - start)/(1e9*86400*365)。pw.udf负责把普通 Python 函数编译进流式计算图并在每条新数据到达时增量调用。
隐含波动率求解(Black-76 模型反推)
定价采用经典 Black(Black-76)模型。脚本中的 compute_price 一次覆盖看涨/看跌:
d1 = (ln(F/K) + σ²·T/2) / (σ·√T)
d2 = d1 - σ·√T
Price = e^(-r·T) · sign · (N(sign·d1)·F − N(sign·d2)·K)
其中 sign = 2·is_call − 1(Call 为 1,Put 为 −1),N(·) 是标准正态分布 CDF,利率 r = 0.043 作为常量参数(见脚本顶部 interest_rate = 0.043)。
由于市场报价给出的是价格,脚本需要反解出该价格对应的隐含波动率:令 BlackPrice(σ) − 期权中间价 = 0,用 scipy.optimize.root_scalar 在 [x0=0.0001, x1=0.8] 区间内求根(greeks-replay.py)。若求根不收敛则返回 None,随后用 filter(pw.this.volatility.is_not_none()) 剔除这些无法定价的期权,避免后续出现 NaN。
五大 Greeks 的 UDF 实现
得到 σ 后,再以 @pw.udf 形式计算 d1/d2 与各 Greeks。源码中五个计算函数的公式(除以 100 表示对 1% 波动的敏感度,theta 除以 252 近似换算为"每个交易日"的时间衰减):
- Delta:Call 为
e^(−rT)·N(d1),Put 为−e^(−rT)·N(−d1); - Gamma:
e^(−rT)·φ(d1) / (F·σ·√T)(Call/Put 相同),其中φ为标准正态密度; - Theta:
[−F·σ·φ(d1)/(2·√T) − r·K·e^(−rT)·N(d2) + r·F·e^(−rT)·N(d1)] / 252(Call),Put 为对称形式/ 252; - Vega:
F·φ(d1)·√T·e^(−rT) / 100; - Rho:Call 为
−T·e^(−rT)·(F·N(d1) − K·N(d2)) / 100,Put 取对称形式。
最终通过一次 select 保留 ts_recv、instrument_id 与五个 Greeks 字段,组成 table_greeks(greeks-replay.py)。
对外输出:注册表 + 运行引擎
port = int(os.environ.get("PORT"))
streamlit_ux.send_table_to_web(port, table_greeks, "table_greeks")
pw.run()
send_table_to_web 会把表注册为别名 table_greeks 并启动查询服务,最后 pw.run() 驱动引擎持续消费与计算。值得一提的是 Databento 返回的是历史区间数据,脚本通过 time.sleep(time_between_updates)(0.05 秒)人为节流,模拟行情逐步到达的效果——代码注释强调:"This is the only difference between static and replay."
主程序二:greeks-static.py——用本地 CSV 离线验证同一套逻辑
greeks-static.py 与 replay 版本共享几乎完全相同的计算链路(过滤 → 聚合中间价 → join → 取远期价格 → 到期时间 → 隐含波动率 → Greeks UDF → 输出 table_greeks),唯一区别在数据入口:
table_es = pw.io.csv.read(definitions_path, schema=DefinitionInputSchema, mode="static")
table_mbp1 = pw.io.csv.read(options_path, schema=OptionInputSchema, mode="static")
即用 pw.io.csv.read(..., mode="static") 一次性读入 data/definition.csv 与 data/options.csv,且不再有 time.sleep 的节流。这种"一个真实数据源、一个 CSV 等价物"的组织方式,非常便于在没有 Databento 账号或需要复现问题时做离线验证:先跑静态脚本确认算法正确,再切到 replay 脚本验证实时链路。
把结果可视化:Streamlit 仪表盘与表格查询服务
启动 UX
无论运行哪个计算脚本,都通过同一命令启动界面:
streamlit run streamlit_ux.py
服务默认监听 http://localhost:8501。刚启动而计算脚本尚未运行/推流时,页面可能因拿不到数据而显示错误提示——这是符合预期的(README 原话:Streamlit 期待被推流的数据,但此时尚未开始)。
计算脚本与 UI 的通信机制
两者并不直接共用进程,而是通过 HTTP 解耦,这条链路由 streamlit_ux.py 与 querying.py 共同实现:
- 计算脚本调用
send_table_to_web(port, table, alias); - 其中
register_table(实现见 streamlit_ux.py)转发到 querying.py 的register_table,其内部用pw.io.subscribe(table, on_change=update)订阅表的增量变更:is_addition=True时把行写入内存字典,否则删除对应键,从而维护一份"当前最新结果"; run_with_querying把pw.run()放在后台线程,主线程启动 uvicorn 承载 FastAPI;/get_table?alias=table_greeks端点把注册表内容序列化为 JSON 返回。若PATHWAY_PROCESS_ID环境变量非 "0"(即已由外部引擎托管),则只运行计算而不重复启动 Web 服务(querying.py);streamlit_ux.py主程序里定义DataModes(STATIC/REPLAY/LIVE)枚举,UI 用带重试的requests.Session轮询http://localhost:{port}/get_table?alias=table_greeks,把响应 JSON 转成 DataFrame、以instrument_id为索引、将ts_recv由纳秒转为可读时间后展示。当前示例中模式固定为STATIC,仅用于界面文案提示,切换脚本运行的是 replay/static 两个计算程序而非该枚举。
仪表盘的实际效果如下(来源:Streamlit.png),表格中可以看到每个 instrument_id 对应的 delta、gamma、theta、vega、rho 与接收时间戳:
运行计算脚本并验证
启动计算脚本:
python greeks-replay.py
等待几秒让脚本完成初始化与数据拉取后,行情便会开始流入。两种验证方式:
- 刷新仪表盘页面(
http://localhost:8501),看到类似上图的表格即代表链路打通;由于处于replay回放模式,数据是带延迟逐步到达的,多次刷新可以看到表格随时间不断增长; - 直接请求查询接口:访问
http://localhost:16001/get_table?alias=table_greeks,或在终端执行curl 'http://localhost:16001/get_table?alias=table_greeks',若返回包含 Greeks 数值的 JSON,说明 Pathway 计算端与查询服务均工作正常。其中 16001 即.env中PORT的值。
运行要点与工程启示
- 依赖外部服务的验证顺序:先跑
python greeks-static.py+ Streamlit 验证算法与展示逻辑(无需 API Key 也能看到完整表格),再运行python greeks-replay.py验证 Databento 接入与回放节奏。两个脚本共享核心计算链,便于快速定位是"算法问题"还是"数据链路问题"; - 计算假设要清晰:示例把无风险利率固定为
r = 0.043、隐含波动率求解区间固定为[0.0001, 0.8],theta 按每年 252 个交易日折算、vega/rho 按 1% 变动标度。这些参数在 greeks-replay.py 与 greeks-static.py 顶部都有注释,实际工程中应根据合约与市场环境调整; - 可复用的"表 → HTTP → 前端"模式:Pathway 表格天然支持订阅式增量输出,本项目用
pw.io.subscribe+ FastAPI 把它桥接成 REST 查询端点,再让 Streamlit 消费。这一模式(实现见 querying.py)同样适用于任意 Pathway 实时结果的前端可视化,而 streamlit_ux.py 中的DataModes枚举也预留了 STATIC / REPLAY / LIVE 三种模式的概念空间,为后续接入真正的实时行情留出了演进方向。
总的来说,这个示例的价值在于它把"金融衍生品定价模型"与"流式计算框架"完整打通:行情数据通过自定义 Connector 进入 Pathway 表,中间价聚合、远期价格引用、隐含波动率反解、五大 Greeks 计算全部以表操作与 UDF 形式在数据流上增量执行,最终结果既可通过 HTTP 接口被任意下游消费,也可直接渲染到 Streamlit 仪表盘——这正是 Pathway 在实时量化风控与交易分析场景中的典型落地范式。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0627
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
