首页
/ Pathway 流式计算期权 Greeks:Option Greeks 示例项目架构、源码与运行全解析

Pathway 流式计算期权 Greeks:Option Greeks 示例项目架构、源码与运行全解析

2026-09-07 17:15:28作者:齐添朝

本示例以 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):

Option Greeks workflow with Pathway Live Data Framework and Databento

示例项目全景与文件组织

本项目源码位于仓库的 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 KeyHTTP 服务端口等参数;
  • 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.pystreamlit_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:标的合约符号(如 ESM4ESZ4);
  • 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    # 期权标识

ConnectorSubjectrun() 中调用 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 的 InstrumentMapinstrument_id 解析回符号;当买/卖价等于 (1<<63) 附近的标记值时视为"未知价格"并跳过;价格统一除以 1e9 归一化(greeks-replay.py)。

数据流上的计算管线

完成数据接入后,全部计算都用 Pathway 表操作表达:

  1. 过滤合约定义:先 filter(pw.this.underlying == front_month_symbol) 只保留近月合约 ESM4 的定义,再过滤出 instrument_class"C""P" 的行;
  2. 提取需要询价的符号集合:通过 reduce(symbol_tuple=pw.reducers.tuple(...)) 聚合出所有期权符号,再经 pw.debug.table_to_pandas(实现见 python/pathway/debug/init.py)拿到符号列表,作为第二段行情请求的 symbols 参数;
  3. 计算期权中间价:对行情表按 raw_symbolgroupby().reduce(),用 pw.reducers.avg 对买价/卖价求均值,得到 option_midprice = (avg(bid) + avg(ask)) / 2
  4. Join 两张表:将"合约定义"与"期权中间价"按 raw_symbol 等值连接(join),产出每个期权最新的中间价;
  5. 取期货价格作为远期价格 F:调用 table_mbp1.ix_ref(front_month_symbol).option_midprice 从行情表中按主键取出 ESM4 期货自身的中间价,作为 Black 模型中的远期价格 Fix_ref 是 Pathway 提供的"按键引用某一行字段"的表操作原语(实现见 python/pathway/internals/table.py);
  6. 计算距离到期时间(年):用 @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)
  • Gammae^(−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
  • VegaF·φ(d1)·√T·e^(−rT) / 100
  • Rho:Call 为 −T·e^(−rT)·(F·N(d1) − K·N(d2)) / 100,Put 取对称形式。

最终通过一次 select 保留 ts_recvinstrument_id 与五个 Greeks 字段,组成 table_greeksgreeks-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.csvdata/options.csv,且不再有 time.sleep 的节流。这种"一个真实数据源、一个 CSV 等价物"的组织方式,非常便于在没有 Databento 账号或需要复现问题时做离线验证:先跑静态脚本确认算法正确,再切到 replay 脚本验证实时链路。

把结果可视化:Streamlit 仪表盘与表格查询服务

启动 UX

无论运行哪个计算脚本,都通过同一命令启动界面:

streamlit run streamlit_ux.py

服务默认监听 http://localhost:8501。刚启动而计算脚本尚未运行/推流时,页面可能因拿不到数据而显示错误提示——这是符合预期的(README 原话:Streamlit 期待被推流的数据,但此时尚未开始)。

计算脚本与 UI 的通信机制

两者并不直接共用进程,而是通过 HTTP 解耦,这条链路由 streamlit_ux.pyquerying.py 共同实现:

  1. 计算脚本调用 send_table_to_web(port, table, alias)
  2. 其中 register_table(实现见 streamlit_ux.py)转发到 querying.pyregister_table,其内部用 pw.io.subscribe(table, on_change=update) 订阅表的增量变更:is_addition=True 时把行写入内存字典,否则删除对应键,从而维护一份"当前最新结果";
  3. run_with_queryingpw.run() 放在后台线程,主线程启动 uvicorn 承载 FastAPI;/get_table?alias=table_greeks 端点把注册表内容序列化为 JSON 返回。若 PATHWAY_PROCESS_ID 环境变量非 "0"(即已由外部引擎托管),则只运行计算而不重复启动 Web 服务(querying.py);
  4. 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 与接收时间戳:

Pathway Option Greeks 计算结果的 Streamlit 仪表盘界面,展示各期权的 delta、gamma、theta、vega、rho 等列

运行计算脚本并验证

启动计算脚本:

python greeks-replay.py

等待几秒让脚本完成初始化与数据拉取后,行情便会开始流入。两种验证方式:

  1. 刷新仪表盘页面http://localhost:8501),看到类似上图的表格即代表链路打通;由于处于 replay 回放模式,数据是带延迟逐步到达的,多次刷新可以看到表格随时间不断增长;
  2. 直接请求查询接口:访问 http://localhost:16001/get_table?alias=table_greeks,或在终端执行 curl 'http://localhost:16001/get_table?alias=table_greeks',若返回包含 Greeks 数值的 JSON,说明 Pathway 计算端与查询服务均工作正常。其中 16001 即 .envPORT 的值。

运行要点与工程启示

  • 依赖外部服务的验证顺序:先跑 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.pygreeks-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 在实时量化风控与交易分析场景中的典型落地范式。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388