Pathway:从 Jupyter 交互式探索到 Kafka 实时布林带监控看板的 Docker 生产部署
本篇基于 Pathway 官方教程 From Jupyter to Deploy,完整复刻一条真实数据科学项目的落地路径:先在 Jupyter 中对静态 CSV 做布林带(Bollinger Bands)交易信号探索,再把同一份代码原样切换为流式数据源打造实时看板,随后接入 Kafka 与 Slack 告警,最终用 Docker Compose 部署为生产级应用。读完并跟做后,你能掌握 Pathway “批流同构(batch and streaming parity)”的核心设计、windowby/reduce 时间窗口聚合、exactly_once_behavior 时序行为配置、Kafka 连接器参数,以及从 Notebook 到可部署 Python 应用的完整转换方法。
教程总览与四阶段路线图
教程要求你在一个高频交易场景中构建带告警能力的实时股票监控看板:当资产价格跌破布林带下轨时提示买入、突破上轨时提示卖出。四部分彼此依赖、层层递进,但也可以独立阅读,各阶段能力矩阵如下:
| 能力 | Part 1: 静态数据探索 | Part 2: 实时看板 | Part 3: Kafka 数据源 | Part 4: 生产部署 |
|---|---|---|---|---|
| Jupyter 中交互式开发 Pathway Live Data Framework | ✅ | |||
| 对静态数据样本做实验 | ✅ | |||
| 流式数据模拟 | ✅ | |||
| Jupyter 内实时可视化 | ✅ | |||
| 告警功能 / 数据输出 | ✅ | |||
| 接入 Kafka 数据源 | ✅ | |||
| 容器化部署 Docker | ✅ | |||
| 将 Jupyter 实时看板嵌入 Web | ✅ |
整条路线体现了构建 Pathway 时的关键设计决策:同一份代码同时适用于有界(batch)与无界(streaming)数据集。Part 1 你只需处理静态数据集;Part 2 直接用相同代码处理数据流,无需操心更新、触发器、算子状态等动态数据集配套问题。但 Pathway 也允许你利用数据的时序特性做扩展:缓冲早到的数据、过滤迟到数据、通过丢弃旧数据清理内存。Part 3 切换到 Kafka 数据源并增加告警,是可部署版本的最后一步;Part 4 把 Part 3 的代码组装成生产级看板。
对应的完整可运行示例都在 examples/projects/from_jupyter_to_deploy 目录下:
- part1_jupyter_exploration.ipynb
- part2_interactive_dashboard.ipynb
- part3_kafka_and_alerts.ipynb
- part3_kafka_data_streamer.ipynb
- part4_deployment 目录(Docker 部署包)
环境准备
Part 1~3 全部在 Jupyter Notebook 中完成,可以本地执行(需要先安装 pathway 并遵循官方 Getting Started 指南),也可以在 Google Colab 云端执行。Part 4 则把项目从 Notebook 转为可部署的 Python 包,需要在 Linux 或 Mac 上安装 Docker。
一个已知的兼容性坑:如果你打算在 VSCode 里运行这些教程 Notebook,由于 jupyter_bokeh 与特定版本 jupyterlab 存在兼容性问题,需要固定 jupyterlab 版本,具体规避方案可参考教程正文中指向的 issue 讨论。
Part 1:在 Jupyter 中探索静态数据
假设你拿到了一份 CSV 格式的历史行情数据,目标是在 Jupyter 中做静态分析——而最关键的收益是:这部分探索中写出的 Pathway 代码在 Part 2 直接就能用于流式数据。
安装依赖与下载数据
在 Notebook 代码单元中安装依赖并拉取数据集:
!pip install pathway
!wget -nc https://gist.githubusercontent.com/janchorowski/e351af72ecd8d206a34763a428826ab7/raw/ticker.csv
加载数据:schema 与 static 模式
用 pw.io.csv.read 把 CSV 载入 Pathway Live Data Framework 的表。这里有一个理解框架行为的关键点:在定义计算逻辑时,数据往往还不存在。因此表示计算的图是带类型的(typed),保证未来流入的数据能被正确处理——pw.io.csv.read 必须提供一个 schema 作为数据结构蓝图。
为了省事,Pathway 提供了 schema_from_csv,直接从某个静态 CSV 文件推断出 schema(其 Python 侧入口位于 python/pathway/internals/schema.py):
import datetime
import pathway as pw
fname = "ticker.csv"
schema = pw.schema_from_csv(fname)
data = pw.io.csv.read(fname, schema=schema, mode="static")
data
最后一行的空 data 语句会打印出表内容。其原理是:Pathway 临时构建并执行一条读取数据并捕获进表的 dataflow。之所以可行,是因为 mode="static" 告诉文件读取器“读完后即退出”——数据预先可知时,Pathway 允许按需预览结果。当数据源改为未来的流式数据时会发生什么,是 Part 2 的主题。
检查 data 表可发现时间戳是整数,需要解析。使用 dt 命名空间下的 utc_from_timestamp:
data = data.with_columns(t=data.t.dt.utc_from_timestamp(unit="ms"))
data
它会返回 pw.DateTimeUtc 类型的通用时间,输出中时间戳即被正确解析并格式化。
用 Bokeh 绘制 VWAP 曲线
从 Live Data Framework 表出图需要外部可视化库。站在流式数据的视角,该库必须支持“随着 Pathway 处理数据持续更新图表”,本教程选择 Bokeh。绘图函数只接收一个 Bokeh ColumnDataSource 参数 src,Pathway 会自动用表内容填充它。以下示例绘制 vwap(成交量加权平均价)的历史曲线:
import bokeh.plotting
def vwap_history(src):
fig = bokeh.plotting.figure(
height=400, width=600,
title="Volume-weighted average price",
x_axis_type="datetime"
)
fig.line("t", "vwap", source=src)
return fig
data.plot(vwap_history, sorting_col="t")
调用 plot 方法并指定 sorting_col="t",即按 t 列排序后作图。
设计算法:布林带统计
布林带方法基于上、下两条边界分别触发卖出与买入告警。边界以一条移动平均为中心,其半径为移动平均标准差的倍数。本例移动平均取 20 分钟周期,半径取 2 倍标准差。
要计算这些统计量,需要按滑动 20 分钟窗口对数据分组,再归约出移动平均 vwap、标准差 vwstd 以及布林带上下轨 bollinger_upper / bollinger_lower。这一分组操作由 windowby 加 reduce 两步完成(窗口机制可进一步查阅 temporal-data 文档目录):
minute_20_stats = (
data
.windowby(
pw.this.t,
window=pw.temporal.sliding(
hop=datetime.timedelta(minutes=1),
duration=datetime.timedelta(minutes=20)
),
instance=pw.this.ticker
)
.reduce(
ticker=pw.this._pw_instance,
t=pw.this._pw_window_end,
volume=pw.reducers.sum(pw.this.volume),
transact_total=pw.reducers.sum(pw.this.volume * pw.this.vwap),
transact_total2=pw.reducers.sum(pw.this.volume * pw.this.vwap**2)
)
.with_columns(
vwap=pw.this.transact_total / pw.this.volume
)
.with_columns(
vwstd=(pw.this.transact_total2 / pw.this.volume - pw.this.vwap**2)**0.5
).with_columns(
bollinger_upper=pw.this.vwap + 2 * pw.this.vwstd,
bollinger_lower=pw.this.vwap - 2 * pw.this.vwstd
)
)
minute_20_stats
参数细节值得注意:
hop是窗口滑动步长,本例为 1 分钟,可按需调整;instance=pw.this.ticker指定独立分析实例——不同股票代码的成交绝不能混在一起统计;_pw_instance列内容即来自该instance参数,_pw_window_end给出窗口右端点,二者在reduce中直接作为输出列。
此外,为了判断买卖时点,还需要即时价格演化。样本代码用 1 分钟**翻滚窗口(tumbling windows,滑动窗口的特例)**计算分钟均值:
minute_1_stats = (
data
.windowby(
pw.this.t,
window=pw.temporal.tumbling(datetime.timedelta(minutes=1)),
instance=pw.this.ticker
)
.reduce(
ticker=pw.this._pw_instance,
t=pw.this._pw_window_end,
volume=pw.reducers.sum(pw.this.volume),
transact_total=pw.reducers.sum(pw.this.volume * pw.this.vwap)
)
.with_columns(
vwap=pw.this.transact_total / pw.this.volume
)
)
minute_1_stats
两张统计表 join 之后,即可自动检测“交易量为何时触发告警”。借助 Pathway 的 if_else,is_alert 列被直接翻译成 buy / hodl / sell 的决策工具:
joint_stats = (
minute_1_stats.join(
minute_20_stats, pw.left.t == pw.right.t, pw.left.ticker == pw.right.ticker
)
.select(
*pw.left,
bollinger_lower=pw.right.bollinger_lower,
bollinger_upper=pw.right.bollinger_upper
)
.with_columns(
is_alert=(
(pw.this.volume > 10000)
& (
(pw.this.vwap > pw.this.bollinger_upper)
| (pw.this.vwap < pw.this.bollinger_lower)
)
)
)
.with_columns(
action=pw.if_else(
pw.this.is_alert,
pw.if_else(pw.this.vwap > pw.this.bollinger_upper, "sell", "buy"),
"hodl"
)
)
)
joint_stats
告警过滤则很容易转化为一张可直接嵌入最终看板的表:
alerts = (
joint_stats
.filter(pw.this.is_alert)
.select(pw.this.ticker, pw.this.t, pw.this.vwap, pw.this.action)
)
alerts
绘制布林带图表
joint_stats 已包含全部所需信息:vwap 画瞬时价格,bollinger_lower 与 bollinger_upper 画布林带,Pathway 触发的告警以散点叠加,得到一个即用型决策组件:
import bokeh.models
def stats_plotter(src):
actions=["buy", "sell", "hodl"]
color_map = bokeh.models.CategoricalColorMapper(
factors=actions,
palette=("#00ff00", "#ff0000", "#00000000")
)
fig = bokeh.plotting.figure(
height=400, width=600,
title="20 minutes Bollinger bands with last 1 minute average",
x_axis_type="datetime"
)
fig.line("t", "vwap", source=src)
fig.line("t", "bollinger_lower", source=src, line_alpha=0.3)
fig.line("t", "bollinger_upper", source=src, line_alpha=0.3)
fig.varea(
x="t",
y1="bollinger_lower",
y2="bollinger_upper",
fill_alpha=0.3,
fill_color="gray",
source=src,
)
fig.scatter(
"t", "vwap",
size=10, marker="circle",
color={"field": "action", "transform": color_map},
source=src
)
return fig
joint_stats.plot(stats_plotter, sorting_col="t")
从图上可以看到该策略通常以低于卖出的价位决定买入——静态数据探索到此结束:你已加载静态数据、用 Live Data Framework 的转换工具发现数据中的模式并生成了有说服力的可视化。
Part 2:从静态探索到实时看板原型
从 Part 1 的 Notebook 出发(复制为新文件或基于 part1_jupyter_exploration.ipynb),把静态数据源换成流式源,复用 Part 1 全部计算代码,做出交互更新的看板。
切换为流式数据:pw.demo.replay_csv
pw.demo.replay_csv 是从静态 CSV 源切换到同数据流式源的最简单方式。它与 pw.io.csv.read 使用相同 schema,额外提供 input_rate 参数控制每秒读取的行数:
# data = pw.io.csv.read(fname, schema=schema, mode="static")
data = pw.demo.replay_csv(fname, schema=schema, input_rate=1000)
# data
从源码看,replay_csv 的实现 内部会按 autocommit_ms = int(1000.0 / input_rate) 推导自动提交周期,并用一个逐行 sleep(1.0 / input_rate) 的 ConnectorSubject 把 CSV 行以 JSON 形式喂入 pw.io.python.read,再把各列从字符串 cast_to_types 回 schema 声明类型——也就是说 input_rate 同时决定了回放节奏与提交粒度(默认值为 1.0 行/秒)。
由于回放沿用最初从输入推断的 schema,时间戳解析需要再做一次:
data = data.with_columns(t=data.t.dt.utc_from_timestamp(unit="ms"))
因为目标变成看板,不再需要展示中间表:删除数据加载单元末尾的空 data 语句,并把其它单元末尾的打印语句注释掉即可。之后,Part 1 为静态数据设计的代码原样生效。
运行看板
用 Panel 把布林带图表与 alerts 表横向排列,show 方法支持隐藏 Pathway 的 id 列、按时间戳排序让最新告警置顶:
import panel as pn
pn.Row(
joint_stats.plot(stats_plotter, sorting_col='t'),
alerts.show(include_id=False, sorters=[{"field": "t", "dir": "desc"}])
)
此时两个组件会显示 Streaming mode(Part 1 里是 Static preview)。最后一步是启动流式处理:
pw.run()
观察到的两个行为特征
- 乱序到达与结果回改:数据点会乱序到达,Pathway 会更新旧结果。一个可见副作用是某些“临时告警”在更多数据到来后消失。这本身就是期望行为,但如果你希望系统稍作等待再输出、减少结果频繁变化,可以通过下面的
behavior配置。 - 无界内存:由于 Pathway 保证输出未来任意时刻都可能被更新,它必须保存(可能是聚合后的)全部历史事件信息,长数据流可能需要无界空间。反过来,你可以向 Pathway“承诺”某些数据点不会再变化,系统便允许释放相关内存;当承诺被打破时,系统被允许忽略这种迟到更新。
这两类调整略微偏离了批流同构——静态模式下所有数据同时到达,本就无需延迟或丢弃——但偏离很容易通过 behavior 配置,从而让代码受益于数据流的时序特性。
为窗口定义时序行为(behavior)
窗口何时产出结果?下游看到的是不完整窗口上的临时聚合,还是完整窗口的最终聚合?系统为迟到数据等多久?这些都由窗口的 behavior 控制。把 behavior 设为 exactly_once_behavior() 后,Pathway 会在窗口收集完全部数据时才首次产出结果,并忽略此后落入该窗口的所有迟到更新;同时利用“某些窗口不会再变化”的信息回收内存。
minute_20_stats = (
data
.windowby(
pw.this.t,
window=pw.temporal.sliding(
hop=datetime.timedelta(minutes=1),
duration=datetime.timedelta(minutes=20)
),
# Wait until the window collected all data before producing a result
behavior=pw.temporal.exactly_once_behavior(),
instance=pw.this.ticker
)
.reduce(
ticker=pw.this._pw_instance,
t=pw.this._pw_window_end,
volume=pw.reducers.sum(pw.this.volume),
transact_total=pw.reducers.sum(pw.this.volume * pw.this.vwap),
transact_total2=pw.reducers.sum(pw.this.volume * pw.this.vwap**2)
)
.with_columns(
vwap=pw.this.transact_total / pw.this.volume
)
.with_columns(
vwstd=(pw.this.transact_total2 / pw.this.volume - pw.this.vwap**2)**0.5
).with_columns(
bollinger_upper=pw.this.vwap + 2 * pw.this.vwstd,
bollinger_lower=pw.this.vwap - 2 * pw.this.vwstd
)
)
minute_1_stats = (
data.windowby(
pw.this.t,
window=pw.temporal.tumbling(datetime.timedelta(minutes=1)),
behavior=pw.temporal.exactly_once_behavior(),
instance=pw.this.ticker,
)
.reduce(
ticker=pw.this._pw_instance,
t=pw.this._pw_window_end,
volume=pw.reducers.sum(pw.this.volume),
transact_total=pw.reducers.sum(pw.this.volume * pw.this.vwap),
)
.with_columns(vwap=pw.this.transact_total / pw.this.volume)
)
从源码看,行为体系定义在 python/pathway/stdlib/temporal/temporal_behavior.py:exactly_once_behavior(shift) 是窗口与时间连接行为的一个便捷工厂,底层是 CommonBehavior 数据类(同文件 L21),包含 delay(延迟首输/加入时间)、cutoff(迟到数据截断)与 keep_results(是否保留结果以支持后续更新)三个基本配置维度。每个时序算子以自己的“当前时间”(输入到达的最大时间)为基准,决定哪些输入或输出被延迟、被忽略,以及何时可以删除不可能再与未来输入交互的内部状态。
加上这些改动后,看板中的图表将呈现“先缓冲、后一次性刷新”的稳定行为,临时闪烁的告警也随之消失。
Part 3:Kafka 集成与告警转发
从 Part 2 的看板出发(复制 part2_interactive_dashboard.ipynb 内容),接入生产级 Kafka 数据源,并把告警推送到 Slack。
准备 Kafka 环境
需要一个正在运行的 Kafka 实例,最容易的两种方式:
- 运行 Kafka Docker 镜像:官方或第三方镜像均可;部分 Kafka 配置依赖 ZooKeeper 管理元数据,除非改用 KRaft 模式(免 ZooKeeper)。可以用 Docker Compose 把两个服务一起跑起来(Part 4 的 compose 文件即采用 ZK 模式);
- 使用托管 Kafka 服务:提供托管实例的云服务多有免费层或试用期。
准备就绪后,创建一个名为 ticker 的 topic。
向 Kafka 写消息:数据流生产者
新建一个辅助 Notebook,把 CSV 数据读入(同 Part 2)再用 pw.io.kafka.write 写入 Kafka topic。先获取 CSV 文件:
!wget -nc https://gist.githubusercontent.com/janchorowski/e351af72ecd8d206a34763a428826ab7/raw/ticker.csv
前两节 Notebook 直接推断 schema,这里改为生成显式的 schema 类以强制约束 Kafka 消息结构:
fname = "ticker.csv"
schema = pw.schema_from_csv(fname)
print(schema.generate_class(class_name="DataSchema"))
然后写出数据流:
# The schema definition is autogenerated
class DataSchema(pw.Schema):
ticker: str
open: float
high: float
low: float
close: float
volume: float
vwap: float
t: int
transactions: int
otc: str
data = pw.demo.replay_csv(fname, schema=DataSchema, input_rate=1000)
rdkafka_producer_settings = {
"bootstrap.servers": "KAFKA_ENDPOINT:9092",
"security.protocol": "sasl_ssl",
"sasl.mechanism": "SCRAM-SHA-256",
"sasl.username": "KAFKA_USERNAME",
"sasl.password": "KAFKA_PASSWORD"
}
pw.io.kafka.write(data, rdkafka_producer_settings, topic_name="ticker")
最后照常用 pw.run() 启动流水线;如果使用 Kafka 云服务,可在控制台确认消息已到达 topic。仓库中该步骤的最终脚本见 part3_kafka_data_streamer.ipynb,其导出形态 kafka-data-streamer.py 中即 pw.io.kafka.write(data, rdkafka_producer_settings, topic_name="ticker")。
从 Kafka 读消息:流式数据源
Pathway 用 pw.io.kafka.read 消费 topic,data 表的填充方式变为:
# Please fill in KAFKA_ENDPOINT, KAFKA_USERNAME, and KAFKA_PASSWORD from your
# cluster configuration.
# Message read status is tracked by consumer group - resetting to a new name
# will cause the program to read messages from the start of the topic.
rdkafka_consumer_settings = {
"bootstrap.servers": "KAFKA_ENDPOINT:9092",
"security.protocol": "sasl_ssl",
"sasl.mechanism": "SCRAM-SHA-256",
"sasl.username": "KAFKA_USERNAME",
"sasl.password": "KAFKA_PASSWORD",
"group.id": "kafka-group-0",
"auto.offset.reset": "earliest"
}
# The schema definition is autogenerated
class DataSchema(pw.Schema):
ticker: str
open: float
high: float
low: float
close: float
volume: float
vwap: float
t: int
transactions: int
otc: str
data = pw.io.kafka.read(
rdkafka_consumer_settings,
topic="ticker",
format="json",
schema=DataSchema
)
配置要点:group.id 标识消费组,消费进度按组跟踪——换成新组名会让程序从 topic 起点重读;auto.offset.reset: earliest 决定新组从最早消息开始;format="json" 指定消息体格式并用 DataSchema 约束解析。此后算法代码完全不变,但数据已来自生产级 Kafka topic——无感复用此前各部分开发的代码,正是 Live Data Framework 的核心优势之一。
告警转发到 Slack
先在 Slack 平台创建应用、安装到工作区并取得 token,再获取目标频道的 ID。然后把回调挂到 alerts 表上:用 pw.io.subscribe 订阅表的变更(Python 侧入口见 python/pathway/io/_subscribe.py),本例只对行新增(row addition)做出反应,每次新增用 requests 调用 Slack API:
import requests
slack_alert_channel_id = "SLACK_CHANNEL_ID"
slack_alert_token = "SLACK_TOKEN"
def send_slack_alert(key, row, time, is_addition):
if not is_addition:
return
alert_message = f'Please {row["action"]} {row["ticker"]}'
print(f'Sending alert "{alert_message}"')
requests.post(
"https://slack.com/api/chat.postMessage",
data="text={}&channel={}".format(alert_message, slack_alert_channel_id),
headers={
"Authorization": "Bearer {}".format(slack_alert_token),
"Content-Type": "application/x-www-form-urlencoded",
},
).raise_for_status()
pw.io.subscribe(alerts, send_slack_alert)
回调签名 (key, row, time, is_addition) 中,is_addition 区分“新增”与“更新”;这正是 exactly_once_behavior 的价值所在——窗口结果只产出一次,回调只在真正的告警发生时触发一次。
在 Notebook 中联调
要测试经真实 Kafka 连接的实时看板与告警,需要重置 Kafka topic 并同时运行两个 Notebook:
- 删除
tickertopic; - 重新创建空的
tickertopic; - 启动实时看板 Notebook;
- 启动数据流生产者 Notebook。
Part 4:从 Jupyter 到独立部署
把数据处理代码与看板从 Notebook 抽离为独立 Python 应用,并容器化。完整成品见 part4_deployment 目录。
导出数据流 Notebook
把 Part 3 的数据流 Notebook导出为 kafka_data_streamer.py(JupyterLab 中 File -> Save and Export Notebook as... -> Executable Script),并注释掉文件开头安装 Pathway、下载数据集的 shell 命令(导出后它们以 get_ipython() 开头)。同时准备前面各步用到的 CSV 数据文件。
导出看板 Notebook
同法把实时看板 Notebook 导出为 dashboard.py,注释掉其中所有 shell 命令。由于不再依赖 Jupyter 展示可视化,需要把末尾的 pw.run() 替换为:
viz_thread = viz.show(threaded=True, port=8080)
try:
pw.run(monitoring_level=pw.MonitoringLevel.ALL)
finally:
viz_thread.stop()
这段代码在独立线程中于 localhost:8080 启动 Web 服务器,pw.run 结束时确保服务器线程被停止;仓库中 dashboard.py 的对应实现 与此完全一致,monitoring_level=pw.MonitoringLevel.ALL 则开启全量运行指标。Web 服务器的更多配置可查阅 Panel 官方文档。
顺带一提:仓库中的 dashboard.py 还包含一行 pw.set_license_key(...),默认值是带遥测的 demo key,注释掉该行即可使用 Community 版,这与教程正文的部署步骤互不冲突。
准备 Docker 配置
创建 Dockerfile 与 docker-compose.yml,让 Docker 同时运行 Kafka、Dashboard 与 Data Streamer 三个应用。
运行 Pathway 应用的 Dockerfile 只负责安装依赖并拷贝文件:
FROM python:3.11
RUN pip install pathway
COPY . .
deployment 编排文件 先配置 Kafka 与 ZooKeeper:
version: '3'
services:
zookeeper:
image: confluentinc/cp-zookeeper:5.5.3
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:5.5.3
depends_on:
- zookeeper
environment:
KAFKA_AUTO_CREATE_TOPICS: true
KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
KAFKA_ADVERTISED_HOST_NAME: kafka
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
KAFKA_BROKER_ID: 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_JMX_PORT: 9991
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
CONFLUENT_SUPPORT_METRICS_ENABLE: false
command: sh -c "((sleep 15 && kafka-topics --create --zookeeper zookeeper:2181 --replication-factor 1 --partitions 1 --topic tickers)&) && /etc/confluent/docker/run "
command 一行在容器启动 15 秒后自动创建 tickers topic(后台执行),避免手工干预。接着添加 dashboard.py 与 kafka-data-streamer.py 两个容器;看板容器需要端口转发以便访问 Web 服务器,sleep 10 确保进程在 Kafka 就绪后启动:
dashboard:
build:
context: ./
depends_on:
- kafka
ports:
- 8080:8080
command: sh -c "sleep 10 && python dashboard.py"
data-streamer:
build:
context: ./
depends_on:
- kafka
command: sh -c "sleep 10 && python kafka-data-streamer.py"
由于 Compose 环境中的 Kafka 未启用认证,还要把两个脚本里的 rdkafka 配置改回明文协议。dashboard.py:
rdkafka_consumer_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "kafka-group-0",
"auto.offset.reset": "earliest",
}
rdkafka_producer_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
}
运行看板
构建并启动全部容器:
docker compose -f "docker-compose.yml" build
docker compose -f "docker-compose.yml" up
容器运行后,终端会打印各服务日志,包括两个 Pathway 应用的输出。打开浏览器访问 localhost:8080 即可看到实时更新的布林带看板,同时告警消息推送到 Slack——至此完成了从“静态数据探索”到“生产级实时流应用”的完整闭环。
小结
本教程演示了一个完整的 Pathway Live Data Framework 数据科学项目:
- Part 1 在 Jupyter 中对静态 CSV 完成布林带算法探索,
mode="static"的按需预览让你先验证逻辑; - Part 2 仅把数据源换成
pw.demo.replay_csv,同一份windowby/reduce/join代码直接跑成流式看板,exactly_once_behavior()让窗口结果稳定且不回刷; - Part 3 用
pw.io.kafka.read/write接入生产级 Kafka 数据源,pw.io.subscribe把告警推送到 Slack; - Part 4 将 Notebook 导出为独立脚本,用 Panel 的
show(threaded=True)嵌入 Web 服务,Docker Compose 一键拉起 Kafka + 看板 + 数据流三个容器。
“批流同构”是贯穿始终的主线:批处理代码天然可升级为流式计算,而时序行为(behavior)与订阅式输出(subscribe)则是在不改动核心算法的前提下补足生产能力的两个正交扩展点。
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 StartedRust0623
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
