首页
/ 基于 Pathway 与 Databento 的期权 Greeks 实时计算指南

基于 Pathway 与 Databento 的期权 Greeks 实时计算指南

2026-09-07 13:32:07作者:凤尚柏Louis

本篇技术指南围绕仓库中“使用 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

总体工作流

Option Greeks 工作流示意图

整个流水线可归纳为:读取合约定义与盘口数据 → 过滤出与目标期货(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_KEYPORT),实际运行前需填入你自己的密钥。

快速开始:完整运行步骤

1. 安装依赖

在原文档与仓库中,运行两个脚本所需的第三方库统一写在 requirements.txt 中,包含 databentopandasscipypathwaypython-dotenvfastapipydanticstreamlituvicorn。安装命令:

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

等待数秒完成启动后,数据即开始以流的方式注入。验证是否成功的两种方式(与原文档一致):

  1. 刷新 Streamlit 页面:看到类似下图的表格即表示成功。由于处于 replay 模式,数据是带延迟地逐条回放,反复刷新可以看到表格在实时增长:

    Streamlit 中的 Option Greeks 结果表格

  2. 直接查询 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.csvdata/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_timedata_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)。

回放脚本中有两个值得注意的数据清洗点:

  1. 合约定义价格需 /1e9 归一化(Databento 传输中使用 1e-9 为单位)。
  2. mbp-1 盘口数据中,买/卖价为 INT64_MAX 附近的占位值代表“未知/无效价格”,代码用 levels[0].bid_px > (1 << 63) - 10 判断并跳过这些记录(greeks-replay.py),否则会把约 922 亿的错误价格带入计算。
  3. 回放与静态的唯一本质差异是回放脚本在每条行情后调用 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.0001x1=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

先由 σ 计算中间量 d1d2,随后五个 @pw.udf 依次完成指标计算。以代码为准,其数学形式为:

  • Delta:Call 为 e^{−rT}·Φ(d1),Put 为 −e^{−rT}·Φ(−d1)
  • Gammae^{−rT}·φ(d1) / (F·σ·√T)(Call/Put 相同);
  • Theta:含 term = −F·σ·φ(d1) / (2√T) 的表达式,Call/Put 分开计算,并统一除以 252(把“年化时间衰减”折算为“每日”,对应一年约 252 个交易日);
  • VegaF·φ(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_idts_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_webstreamlit_ux.py)实际做了两件事:

  1. 调用 querying.register_table(table, alias):底层通过 pw.io.subscribe(self, on_change=update) 订阅表变更(新增行写入字典、删除行移除),把实时结果缓存在内存中(querying.py);
  2. 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.readpw.io.python.readConnectorSubject,还是其他内置连接器),后续的过滤、join、pw.udf 数值计算、ix_ref 引用与 pw.io.subscribe 输出都可以原样复用。将其用于自己的期权风控场景时,通常只需修改四处:

  1. .env 中的 API_KEYPORT
  2. 顶部参数区(数据集、符号集、主力合约、利率与时间窗);
  3. 若标的不是期货期权,需将 compute_price 及 Greeks 公式替换为对应的定价模型(如标的为股票/指数的 Black-Scholes 或二叉树模型);
  4. 将 Streamlit 页面请求的别名与端口对齐 querying.py 中注册的表名。

原文档与仓库中的完整可运行版本分别见 docs/2.developers/7.templates/ETL/_readmes/option-greeks.mdexamples/projects/option-greeks

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