基于 Pathway 与 Databento 的期权 Greeks 实时计算指南
本篇技术指南围绕仓库中“使用 Pathway 与 Databento 计算期权 Greeks”的完整示例展开,说明如何以金融衍生品数据为输入,在 Pathway 流处理框架中完成合约过滤、盘口数据聚合、隐含波动率求解与 Delta/Gamma/Theta/Vega/Rho 五大风险指标计算,并通过 Streamlit 仪表盘实时展示。读者学完后,将掌握 Pathway 的 CSV 静态读取、Python 连接器回放(replay)、pw.udf 数值计算、表连接/索引引用与 Web 查询(querying)接口的完整实战用法。
项目要解决什么问题:期权的 Greeks 是什么
期权(Options)是赋予持有者权利而非义务、在规定期限内按约定价格买入或卖出标的资产的金融衍生品。仅仅知道期权价格是不够的——交易员与风险管理者还需要量化地回答“价格对哪些因素有多敏感”,这组敏感性指标就是 Option Greeks(期权希腊字母),用于衡量价格、标的资产价格、时间、波动率与利率等多种风险因子变动时期权价格的变化。
本示例基于 Pathway 的计算流程计算并输出五个常用 Greeks:
| 指标 | 含义 | 代码中对应的计算函数 |
|---|---|---|
| Delta(Δ) | 期权价格对标的资产价格变动的敏感度 | compute_delta |
| Gamma(Γ) | Delta 对标的资产价格变动的敏感度(Delta 的变化率) | compute_gamma |
| Theta(Θ) | 期权价格对时间流逝的敏感度(时间衰减) | compute_theta |
| Vega(ν) | 期权价格对隐含波动率变动的敏感度 | compute_vega |
| Rho(ρ) | 期权价格对无风险利率变动的敏感度 | compute_rho |
数据链路为:Databento(提供机构级行情数据的服务商)负责行情与合约数据的获取,Pathway 负责持续计算,Streamlit 负责结果展示。本示例完整代码位于 examples/projects/option-greeks,README 说明见 examples/projects/option-greeks/README.md。
总体工作流
整个流水线可归纳为:读取合约定义与盘口数据 → 过滤出与目标期货(ESM4)相关的 Call/Put → 计算各期权盘中价(midprice) → 求得到期时间与远期价格 → 用 Black 模型反解隐含波动率 → 计算五种 Greeks → 通过 Web 接口对外暴露,供 Streamlit 轮询展示。
项目组织与文件职责
原文档给出的工程目录如下(相对项目根目录的实际路径见各注释):
examples/projects/option-greeks/
├── data/
│ ├── definition.csv # 静态脚本使用的期权合约定义数据
│ └── options.csv # 静态脚本使用的盘口(mbp-1)行情数据
├── .env # 存放 Databento API key 与 Web 服务端口等参数
├── greeks-replay.py # 回放模式:从 Databento Historical API 读取行情并计算 Greeks
├── greeks-static.py # 静态模式:读取上述 CSV 文件并计算 Greeks
├── requirements.txt # 运行所需依赖清单
└── streamlit_ux.py # 基于 Streamlit 的结果展示界面
补充说明:
data/中的数据规模(当前仓库版本)为definition.csv约 3253 行、options.csv约 18545 行。静态 CSV 中合约定义列包含ts_recv, raw_symbol, expiration, instrument_class, strike_price, underlying, instrument_id,盘口数据列包含ts_recv, symbol, bid_px_00, ask_px_00, raw_symbol, bid_px, ask_px,与代码中定义的两套 Schema 一一对应。querying.py不属于原文档树形图,但被streamlit_ux.py依赖,负责用 FastAPI 将 Pathway 表注册为可按别名查询的 HTTP 端点。.env在仓库中保存为模板(API_KEY、PORT),实际运行前需填入你自己的密钥。
快速开始:完整运行步骤
1. 安装依赖
在原文档与仓库中,运行两个脚本所需的第三方库统一写在 requirements.txt 中,包含 databento、pandas、scipy、pathway、python-dotenv、fastapi、pydantic、streamlit、uvicorn。安装命令:
pip install -r requirements.txt
2. 配置 Databento API Key
先到 Databento 注册账号并获取 API key(注册可获得免费额度)。然后打开项目目录下的 .env,填入:
API_KEY = "你的_databento_api_key"
PORT = "16001"
其中 API_KEY 用于调用 Databento Historical API,PORT 是后续 Web 查询服务监听的端口。两个脚本中都会通过 load_dotenv() 读取该文件;脚本内部以 os.environ.get("API_KEY")、os.environ.get("PORT") 获取对应变量。注意:若直接使用静态数据(greeks-static.py),不需要真实 API key,但脚本仍会读取 .env 中的 PORT。
3. 启动 Streamlit UX
无论最终跑静态还是回放脚本,展示界面都以相同方式启动:
streamlit run streamlit_ux.py
Streamlit 服务默认运行在 http://localhost:8501。启动后页面应显示 "Option Greeks" 标题与模式说明文字;此时可能看到“等待数据”的错误提示,这是正常现象——因为后端数据流尚未启动(源码中该页面会请求 http://localhost:{port}/get_table?alias=table_greeks,见 streamlit_ux.py)。
4. 运行数据脚本
需要拿到真实行情做回放时,执行:
python greeks-replay.py
等待数秒完成启动后,数据即开始以流的方式注入。验证是否成功的两种方式(与原文档一致):
-
刷新 Streamlit 页面:看到类似下图的表格即表示成功。由于处于
replay模式,数据是带延迟地逐条回放,反复刷新可以看到表格在实时增长: -
直接查询 JSON 端点:浏览器访问
http://localhost:16001/get_table?alias=table_greeks,若返回包含 Greeks 数值的 JSON,说明后端计算链路工作正常。该接口由querying.py中的 FastAPI 端点实现(别名不存在时返回 404,见 querying.py)。
5. 运行静态脚本
如果只想在不依赖 Databento API 的情况下复现整条计算链路,可以运行使用仓库内置 CSV 的静态版本:
python greeks-static.py
该脚本不调用行情服务,直接读取 data/definition.csv 与 data/options.csv,其余计算逻辑与回放脚本完全一致。
源码级拆解:从行情到 Greeks 的五步流水线
两套脚本的计算主链基本相同,这里以逻辑更完整的 greeks-replay.py 为主线,并标注与静态版 greeks-static.py 的差异。
3.1 定义关键参数与输入 Schema
两类脚本使用同一批行情参数(回放脚本见 greeks-replay.py):
db_dataset = "GLBX.MDP3":CME Globex 的 MDP 3.0 行情数据集,示例中选取的是基于 E-mini S&P 500(ES)的期货期权;db_def_schema = "definition"与db_price_schema = "mbp-1":分别指合约定义(静态属性)与一档盘口行情;db_def_symbols = ["ES.OPT"]:父级符号模式下筛选出所有根代码为ES的期权;front_month_symbol = "ESM4":近月主力期货合约代码,用于定位期权标的;interest_rate = 0.043:Black 模型中的无风险利率(示例为常数 4.3%);- 回放脚本额外需要
start_time、data_duration(获取合约定义的窗口)与query_data_duration(获取行情的时间窗),示例为从 2024-04-04 17:00(Us/Central)起取 2 分钟mbp-1数据。
两个输入都通过 pw.Schema 声明列类型:
class DefinitionInputSchema(pw.Schema):
ts_recv: int # 收到数据的时间(ns)
raw_symbol: str # 期权符号
expiration: int # 期权到期时间(ns)
instrument_class: str # 期权类型 C/P/T 等
strike_price: float # 行权价
underlying: str # 第一标的资产符号
instrument_id: int # 期权唯一标识
class OptionInputSchema(pw.Schema):
raw_symbol: str # 期权符号
bid_px: float # 买价
ask_px: float # 卖价
3.2 读取数据:静态 CSV vs 历史回放
- 静态模式使用
pw.io.csv.read(..., mode="static"),一次性读入 CSV(greeks-static.py)。 - 回放模式使用
pw.io.python.read()配一个ConnectorSubject子类,在run()中调用 Databento 的client.timeseries.get_range(...)逐行取出数据,再通过self.next(...)送入 Pathway(greeks-replay.py)。
回放脚本中有两个值得注意的数据清洗点:
- 合约定义价格需
/1e9归一化(Databento 传输中使用 1e-9 为单位)。 mbp-1盘口数据中,买/卖价为INT64_MAX附近的占位值代表“未知/无效价格”,代码用levels[0].bid_px > (1 << 63) - 10判断并跳过这些记录(greeks-replay.py),否则会把约 922 亿的错误价格带入计算。- 回放与静态的唯一本质差异是回放脚本在每条行情后调用
time.sleep(time_between_updates)(0.05 秒/条)以模拟慢速实时流(greeks-replay.py)——这正是页面刷新时表格不断增长的原因。
3.3 合约过滤、盘口聚合与表连接
拿到两张“表”后,先用连续 filter 收缩范围:
table_esm4 = table_esm4.filter(pw.this.underlying == front_month_symbol) # 只留 ES 近月标的
table_esm4 = table_esm4.filter(
(pw.this.instrument_class == "C") | (pw.this.instrument_class == "P") # 只留 Call/Put
)
回放脚本中还需要把过滤后的 raw_symbol 收集成列表(reduce + pw.reducers.tuple),作为下一步向 Databento 请求期权行情价格的符号集。
接着用 groupby(...).reduce(...) 对同一期权在回放窗口内的买/卖价求平均,得到盘口中价(midprice):
table_mbp1 = table_mbp1.groupby(pw.this.raw_symbol).reduce(
raw_symbol=pw.this.raw_symbol,
option_midprice=(pw.reducers.avg(pw.this.bid_px) + pw.reducers.avg(pw.this.ask_px)) / 2,
)
随后把“合约定义”与“行情价格”两张表以 raw_symbol 等值连接起来:
table_prices = table_esm4.join(
table_mbp1, pw.left.raw_symbol == pw.right.raw_symbol
).select(
*pw.left, # 带入合约定义的整行
option_midprice=pw.right.option_midprice,
)
最后用 ix_ref(按 key 引用另一张表的列)取出主力期货 ESM4 的 midprice 作为期权定价所需的远期价格 F:
table_prices = table_prices.with_columns(
future_price=table_mbp1.ix_ref(front_month_symbol).option_midprice
)
3.4 到期时间与 Black 模型定价
期权定价需要“距离到期还有多少年”。原始 expiration 是纳秒时间戳,因此用一个 @pw.udf 装饰的纯 Python 函数把纳秒差折算成年份:
@pw.udf
def compute_time_to_expiration(expiration_time: int) -> float:
return (expiration_time - int(start_time.timestamp() * 1e9)) / (1e9 * 86400 * 365)
本项目针对的是期货期权,因此使用适合以期货价格 F 为标的的 Black 模型(而非经典的 Black-Scholes 公式)。定价函数返回基于正态分布累计分布函数 norm.cdf 的价格:
def compute_price(F, K, T, sigma, r=interest_rate, is_call=True) -> float:
d1 = (math.log(F / K) + (sigma**2 / 2) * T) / (sigma * math.sqrt(T))
d2 = d1 - sigma * math.sqrt(T)
sign = 2 * int(is_call) - 1
return math.exp(-r * T) * sign * (norm.cdf(sign * d1) * F - norm.cdf(sign * d2) * K)
其中 F 为远期价格,K 为行权价,T 为以年计的到期时间,sigma 为波动率,r 为无风险利率。
3.5 用 scipy 反解隐含波动率
模型给定价需要波动率 σ,但行情并不直接提供它。工程上的做法是反解隐含波动率:找到一个 σ,使 Black 模型价格等于市场 midprice,即求方程
BlackPrice(σ) − midprice = 0
的根。代码用 scipy.optimize.root_scalar,以 x0=0.0001、x1=0.8 为初值区间求解;求解失败(not result.converged)时返回 None:
@pw.udf
def compute_volatility(F, K, T, is_call, option_midprice) -> float | None:
result = scipy.optimize.root_scalar(
lambda sigma: option_midprice - compute_price(F=F, K=K, T=T, sigma=sigma, is_call=is_call),
x0=0.0001,
x1=0.8,
)
return result.root if result.converged else None
由于 Pathway 是增量/持续计算引擎,这里把 scipy 的求根逻辑封装为 pw.udf 后,每条新到价的行会按需重新求根。产生 None 的行随后被 filter(pw.this.volatility.is_not_none()) 剔除,保证后续 Greeks 计算只作用于收敛成功的记录(greeks-replay.py)。
3.6 计算五种 Greeks
先由 σ 计算中间量 d1、d2,随后五个 @pw.udf 依次完成指标计算。以代码为准,其数学形式为:
- Delta:Call 为
e^{−rT}·Φ(d1),Put 为−e^{−rT}·Φ(−d1); - Gamma:
e^{−rT}·φ(d1) / (F·σ·√T)(Call/Put 相同); - Theta:含
term = −F·σ·φ(d1) / (2√T)的表达式,Call/Put 分开计算,并统一除以 252(把“年化时间衰减”折算为“每日”,对应一年约 252 个交易日); - Vega:
F·φ(d1)·√T·e^{−rT} / 100,除以 100 对应“波动率变动 1 个百分点”的价格变化; - Rho:Call 为
−T·e^{−rT}·(F·Φ(d1)−K·Φ(d2)) / 100,Put 为对应负向组合再/100,除以 100 对应利率变动 1 个百分点。
(φ 为标准正态概率密度 norm.pdf,Φ 为标准正态累计分布 norm.cdf。)
五种指标函数都以 @pw.udf 声明(见 greeks-replay.py 与静态脚本中完全相同的实现),最后用一个 select 把五个结果与 instrument_id、ts_recv 一起投影为最终输出表 table_greeks:
table_greeks = table_d1d2.select(
ts_recv=pw.this.ts_recv,
instrument_id=pw.this.instrument_id,
delta=compute_delta(...),
gamma=compute_gamma(...),
theta=compute_theta(...),
vega=compute_vega(...),
rho=compute_rho(...),
)
3.7 结果对外暴露:Web 查询与 Streamlit
最终结果表通过 streamlit_ux.py 的辅助函数暴露给外部:
streamlit_ux.send_table_to_web(port, table_greeks, "table_greeks")
pw.run()
send_table_to_web(streamlit_ux.py)实际做了两件事:
- 调用
querying.register_table(table, alias):底层通过pw.io.subscribe(self, on_change=update)订阅表变更(新增行写入字典、删除行移除),把实时结果缓存在内存中(querying.py); - 以
monitoring_level=pw.MonitoringLevel.NONE在独立线程中启动pw.run(),同时用 Uvicorn 启动监听0.0.0.0:port的 FastAPI 服务,注册GET /get_table?alias=...端点,将表内容序列化为 JSON(querying.py)。
Streamlit 界面(streamlit_ux.py)通过 requests(带重试适配器)周期性拉取 http://localhost:{PORT}/get_table?alias=table_greeks,把 JSON 转为 DataFrame,将纳秒级 ts_recv 转为时间并设置 instrument_id 为索引后,用 st.dataframe 渲染表格;请求失败则显示错误信息。图 Streamlit.png 即为运行成功后的真实界面:表格以 instrument_id 为索引,列出每个期权的 delta/gamma/rho/theta/vega 与数据接收时间 ts_recv。
小结:把这个示例迁移到自己的行情源
这个示例的价值在于它演示了一套与数据源解耦的持续计算范式:只要把表数据送进 Pathway(无论是 pw.io.csv.read、pw.io.python.read 的 ConnectorSubject,还是其他内置连接器),后续的过滤、join、pw.udf 数值计算、ix_ref 引用与 pw.io.subscribe 输出都可以原样复用。将其用于自己的期权风控场景时,通常只需修改四处:
.env中的API_KEY与PORT;- 顶部参数区(数据集、符号集、主力合约、利率与时间窗);
- 若标的不是期货期权,需将
compute_price及 Greeks 公式替换为对应的定价模型(如标的为股票/指数的 Black-Scholes 或二叉树模型); - 将 Streamlit 页面请求的别名与端口对齐
querying.py中注册的表名。
原文档与仓库中的完整可运行版本分别见 docs/2.developers/7.templates/ETL/_readmes/option-greeks.md 与 examples/projects/option-greeks。
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
