Pathway 实时日志监控实战:Filebeat + Kafka + Pathway 滑动窗口告警并推送 Slack
实时监控服务器日志(如 Nginx 访问日志)并第一时间发现流量异常,是流式处理最典型的落地场景之一。本文基于开源仓库 pathway 中 examples/projects/realtime-log-monitoring/filebeat-pathway-slack 这套端到端示例,完整讲解如何用 Filebeat 采集日志 → Kafka 中转 → Pathway 实时计算滑动窗口 → Slack 告警 搭建监控链路。读完本文,你将掌握 Pathway 的 Kafka 输入连接器、ISO 8601 时间解析、基于 windowby 的滑动窗口与数据遗忘行为,以及通过 pw.io.subscribe 自定义输出回调的完整工程化写法。
本文对应的关联文档为 filebeat-pathway-slack.md,示例工程完整源码位于 examples/projects/realtime-log-monitoring/filebeat-pathway-slack,同一目录下还提供了配套的教程页 7.realtime-log-monitoring.md 可供交叉阅读。
一、项目要解决什么问题
日志是系统运行状态、性能与错误信息的重要载体,但人工盯日志既耗时又容易漏报。该示例演示的目标是:监控日志流,当单位时间内日志条数异常激增时,立即向 Slack 频道推送告警。具体业务规则为:
- 以每条日志自带的
@timestamp为时间基准; - 只统计"最近 X 秒内"(默认 X=1s)到达的日志;
- 若窗口内日志条数超过 Y(默认 Y=5),则输出
alert=True并触发 Slack 消息。
为了端到端可复现,工程用四个 Docker 容器把整条链路跑起来:Filebeat(日志来源 + 采集)、Kafka 与 ZooKeeper(中转网关)、Pathway(实时处理与告警判定)、再加上 Slack 作为外部告警接收方。这与官方仓库中的配套演示架构一致:Filebeat 直接以 Kafka 为 output,Pathway 订阅同一个 Kafka topic 完成处理,省去了 Logstash 等中间环节,链路更短、响应更快。
二、工程目录与文件职责
示例工程(examples/projects/realtime-log-monitoring/filebeat-pathway-slack)目录结构如下:
filebeat-pathway-slack/
├── filebeat-src/
│ ├── Dockerfile # 基于 filebeat:8.6.1 镜像,拷入配置与造数脚本
│ ├── filebeat.docker.yml # Filebeat 采集与输出配置(输入 /input_stream/*,输出到 Kafka)
│ └── generate_input_stream.sh # 人工生成日志流的脚本(模拟流量突刺)
├── pathway-src/
│ ├── Dockerfile # 基于 python:3.10 安装 pathway 并运行 alerts.py
│ └── alerts.py # Pathway 核心逻辑:读 Kafka → 滑窗统计 → Slack 告警
├── docker-compose.yml # 编排 filebeat / zookeeper / kafka / pathway 四个容器
└── Makefile # build/stop/connect/connect-pathway 快捷命令
其中 pathway-src/alerts.py 是整条业务链路的"大脑";docker-compose.yml 定义容器编排;Makefile 提供一键启停入口。
三、端到端数据流:日志从文件到 Slack 的四步处理
关联文档将日志处理流程概括为 4 个步骤(全部实现在 pathway-src/alerts.py 中):
- 字段抽取:从 Filebeat 生成的 JSON 消息中取出
@timestamp与message两个字段; - 时间解析:把 ISO 8601 格式的字符串时间转换为真正的(可运算的)时间戳;
- 滑动窗口过滤:只保留最近 X 秒(默认 X=1s)内的日志,其中"当前时间"以窗口内最后一条日志的时间戳为参照;
- 阈值判定:若窗口内日志数超过 Y(默认 Y=5),输出
alert=True。
落到 Pathway 代码上,alerts.py 依次完成参数声明、Kafka 读取、时间清洗、滑窗聚合与告警输出五个环节:
import time
from datetime import timedelta
import requests
import pathway as pw
# 如需使用 Pathway Live Data Framework Scale 的进阶功能,可在此填入从官方渠道申请的 license key;
# 使用 Community 版本时把这一行注释掉即可。
pw.set_license_key("demo-license-key-with-telemetry")
alert_threshold = 5
sliding_window_duration = timedelta(seconds=1)
SLACK_ALERT_CHANNEL_ID = "XXX"
SLACK_ALERT_TOKEN = "XXX"
rdkafka_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "0",
"session.timeout.ms": "6000",
}
接下来按行拆解每一段的关键设计。
3.1 用 schema_builder 声明输入结构
Filebeat 写入 Kafka 的 JSON 消息包含很多元数据字段,本示例只关心 @timestamp 和 message 两列,其余字段可忽略。由于列名 @timestamp 含有 @ 特殊字符(无法用常规类属性访问),代码采用**内联 schema(inline schema)**方式定义:
inputSchema = pw.schema_builder(
columns={
"@timestamp": pw.column_definition(dtype=str),
"message": pw.column_definition(dtype=str),
}
)
这里 pw.schema_builder 把两列都声明为 str 类型——因为 Filebeat 的 JSON 中时间字段是字符串,真正的类型转换要留到下一步由 .dt.strptime 完成。
3.2 通过 Kafka 输入连接器读取日志
Pathway 通过 pw.io.kafka.read 提供的 Kafka 输入连接器订阅 logs topic:
log_table = pw.io.kafka.read(
rdkafka_settings,
topic="logs",
format="json",
schema=inputSchema,
autocommit_duration_ms=100,
)
rdkafka_settings是透传给底层 librdkafka 的参数,其中bootstrap.servers指向 docker-compose 内的 Kafka 服务名kafka:9092(容器内互访不能用localhost);topic="logs"与 Filebeat 的output.kafka配置、以及 Kafka 容器启动时自动创建的主题保持一致;format="json"表明消息体为 JSON;autocommit_duration_ms=100控制消费位点自动提交的间隔。
如果读者需要用到带认证的 Kafka(SASL/SSL),可参考同一目录教程 7.realtime-log-monitoring.md 中给出的带 sasl_ssl、SCRAM-SHA-256 等参数的 rdkafka_settings 写法。
3.3 列重命名与 ISO 8601 时间解析
读取到的表需要两处清洗:一是把含 @ 的列名改为 timestamp/log,方便使用 Pathway 的点号(dot notation)访问;二是把 ISO 8601 字符串时间转成可运算的 datetime:
log_table = log_table.select(timestamp=pw.this["@timestamp"], log=pw.this.message)
log_table = log_table.select(
pw.this.log,
timestamp=pw.this.timestamp.dt.strptime("%Y-%m-%dT%H:%M:%S.%fZ"),
)
strptime 的格式串 %Y-%m-%dT%H:%M:%S.%fZ 精确匹配 Filebeat 默认输出的 UTC 时间戳(含毫秒与结尾的 Z),解析后 timestamp 即成为可参与窗口计算的 datetime 列——这是后续 windowby 能按时间正确切窗的前提。
3.4 用 windowby 构建 1 秒滑动窗口
Pathway 通过 windowby 时间窗口 API 完成流式窗口计算,配合数据遗忘(data forgetting)行为,使程序既能基于"最近的日志"判定告警,又能以常量内存长时间运行:
t_sliding_window = log_table.windowby(
log_table.timestamp,
window=pw.temporal.sliding(
hop=timedelta(milliseconds=10), duration=sliding_window_duration
),
behavior=pw.temporal.common_behavior(
cutoff=timedelta(seconds=0.1),
keep_results=False,
),
).reduce(timestamp=pw.this._pw_window_end, count=pw.reducers.count())
关键参数的含义:
| 参数 | 取值 | 作用 |
|---|---|---|
window=pw.temporal.sliding(hop, duration) |
hop=10ms,duration=1s |
每 10 毫秒滑动生成一个长度为 1 秒的滑动窗口,相邻窗口互相重叠,避免"切窗错位"漏检突发流量 |
behavior=pw.temporal.common_behavior(cutoff=0.1s, keep_results=False) |
cutoff 为 0.1 秒,不保留结果 | 窗口晚于 cutoff 即被遗忘,保证告警只看最近数据、内存有上界 |
reduce(timestamp=pw.this._pw_window_end, count=pw.reducers.count()) |
聚合输出窗口结束时间与条数 | count() 归约器统计落入每个窗口的日志条数 |
流式监控场景中滑动窗口通常优于滚动窗口(tumbling window):滚动窗口把数据切分为互不重叠的区间,一旦切割边界不对就会漏掉目标模式;滑动窗口则始终在"当前时刻往前推一段固定时长"上统计,对突发、短促的流量尖峰更敏感。
3.5 阈值判定与告警表
有了每个滑动窗口的 count 后,取所有窗口的最大条数并与阈值比较,得到单行布尔告警表:
t_alert = t_sliding_window.reduce(count=pw.reducers.max(pw.this.count)).select(
alert=pw.this.count >= alert_threshold
)
逻辑是:只要存在任何一个 1 秒滑动窗口内日志条数 ≥ 5,就认为处于告警状态(alert=True)。由于 Pathway 是增量更新的,每当新日志到达,窗口计数、最大条数与最终布尔值都会被自动重算并向下游推送变化(upsert 语义),无需开发者手动维护窗口集合。
3.6 用 pw.io.subscribe 把告警推送到 Slack
Pathway 本身没有开箱即用的 Slack 输出连接器,因此示例使用 subscribe 输出回调 + requests 直接调用 Slack Web API:
def on_alert_event(key, row, time, is_addition):
alert_message = "Alert '{}' changed state to {}".format(
row["alert"],
"ACTIVE" if is_addition else "INACTIVE",
)
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(t_alert, on_alert_event)
回调函数签名为 (key, row, time, is_addition):row["alert"] 是当前布尔值,is_addition 为 True 表示该行进入表(告警状态变为 ACTIVE),为 False 表示被删除(回到 INACTIVE)——正是这种"状态翻转即通知"的机制,避免了对同一条告警反复轰炸。消息通过 Slack 的 chat.postMessage 接口发送到 SLACK_ALERT_CHANNEL_ID,鉴权使用 SLACK_ALERT_TOKEN Bearer Token。
最后,程序先 time.sleep(5) 等待 Kafka 就绪(否则启动阶段会出现连接报错,虽然不影响最终运行),再调用 pw.run() 启动整个计算图;没有 pw.run(),前面定义的所有表都不会真正开始计算:
time.sleep(5)
# 启动计算
pw.run()
四、容器编排:docker-compose 与 Makefile
pathway-src/Dockerfile 在 python:3.10 基础上安装 pathway、requests、python-dateutil,并把 alerts.py 拷入镜像后以 python -u alerts.py 启动:
FROM --platform=linux/x86_64 python:3.10
RUN pip install -U pathway
RUN pip install requests
RUN pip install python-dateutil
COPY ./pathway-src/alerts.py alerts.py
CMD ["python", "-u", "alerts.py"]
(为兼容性起见,示例固定使用 x86_64 的 Linux 基础镜像;其中 requests 供 Slack HTTP 推送使用。)
docker-compose.yml 定义了四个服务,注意服务名 filebeat/pathway 只是 compose 内的服务别名,实际镜像内容由各自的 build.dockerfile 决定:
version: "3.7"
services:
filebeat:
build:
context: .
dockerfile: ./filebeat-src/Dockerfile
zookeeper:
image: confluentinc/cp-zookeeper:5.5.3
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-enterprise-kafka:5.5.3
depends_on: [zookeeper]
environment:
KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
KAFKA_BROKER_ID: 1
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_JMX_PORT: 9991
ports:
- 9092:9092
command: sh -c "((sleep 15 && kafka-topics --create --zookeeper zookeeper:2181 --replication-factor 1 --partitions 1 --topic logs)&) && /etc/confluent/docker/run "
pathway:
build:
context: .
dockerfile: ./pathway-src/Dockerfile
depends_on: [filebeat]
值得注意的三点实现细节:
- topic 自动创建:Kafka 容器的
command在启动后延迟 15 秒,用kafka-topics --create创建名为logs的单分区主题,再拉起 Kafka 主进程。若想换主题名,需要同步修改这一行、alerts.py 里的topic="logs"以及 Filebeat 的output.kafka配置。 - 容器内互联用服务名:Pathway 中
rdkafka_settings的bootstrap.servers必须写kafka:9092(compose 服务名),Filebeat 的 Kafka output host 同样写kafka:9092,而不是localhost。 - 服务间依赖:
kafka依赖zookeeper,pathway依赖filebeat,通过depends_on控制启动顺序,保证连接对象先就绪。
Makefile 把常用操作收敛为四个目标:
build:
docker compose up -d
stop:
docker compose down -v
connect:
docker compose exec filebeat bash
connect-pathway:
docker compose exec pathway bash
make(即make build)以守护方式拉起全部容器;make connect/make connect-pathway分别进入 Filebeat 与 Pathway 容器;make stop用docker compose down -v停止并连带删除卷,保证下次启动环境干净。
五、Filebeat 配置:监控目录并把日志写进 Kafka
filebeat-src/Dockerfile 基于官方 docker.elastic.co/beats/filebeat:8.6.1 镜像,把配置文件与造数脚本拷入容器,并预先创建空日志文件 /input_stream/example.log:
FROM docker.elastic.co/beats/filebeat:8.6.1
COPY ./filebeat-src/filebeat.docker.yml /usr/share/filebeat/filebeat.yml
COPY ./filebeat-src/generate_input_stream.sh /usr/share/filebeat/generate_input_stream.sh
USER root
RUN mkdir /input_stream/
RUN touch /input_stream/example.log
RUN chown root:filebeat /usr/share/filebeat/filebeat.yml
RUN chmod go-w /usr/share/filebeat/filebeat.yml
Filebeat 运行配置 filebeat.docker.yml 是本链路"数据从哪来、到哪去"的关键:
filebeat.inputs:
- type: filestream
id: my-logs
paths:
- /input_stream/*
filebeat.config.modules:
path: /usr/share/filebeat/modules.d/
reload.enable: false
output.kafka:
enabled: true
hosts: ["kafka:9092"]
topic: "logs"
group_id: 1
ssl.enabled: false
- 输入侧使用
filestream类型,监控目录/input_stream/*下的所有日志文件,id: my-logs用于标识该输入流、维持采集状态; - 输出侧直接启用
output.kafka,把每条日志作为 JSON 消息发往 Kafka 的logs主题(group_id: 1为发送端标识,与消费端无关);ssl.enabled: false表明示例不启用认证,生产环境如需加密可参考 Elastic 官方 Kafka output 文档补充配置; - 若要监控真实 Nginx 日志,只需把
paths改成 Nginx 日志所在目录;Filebeat 也提供针对 Nginx 的专用 module,可进一步提升字段解析能力。
六、模拟流量突刺:generate_input_stream.sh
filebeat-src/generate_input_stream.sh 在缺少真实高并发流量的场景下,用一段 bash 脚本模拟"先平稳、后突刺"的日志模式:
#!/bin/bash
src="../../../input_stream/example.log"
sleep 1
for LOOP_ID in {1..100}
do
printf "$LOOP_ID\n" >> $src
sleep 1
done
for LOOP_ID in {101..200}
do
printf "$LOOP_ID\n" >> $src
sleep 0.01
done
- 脚本进入 Filebeat 容器后的实际位置是
/usr/share/filebeat/,相对路径../../../input_stream/example.log最终解析为/input_stream/example.log,正好落在 Filebeat 监控目录下; - 前半段以 每秒 1 条 的速度追加 100 条日志,模拟正常平稳流量;
- 后半段以 每 10 毫秒 1 条(约每秒 100 条)追加 100 条日志,模拟突发流量尖峰——这一段的写入速率远超 1 秒 5 条的告警阈值,从而稳定触发 Slack 告警。
七、安装、启动与运行验证
7.1 安装前提与配置
安装前需要确认本机已具备 Docker 与 docker compose 环境。整个示例不需要在宿主机安装 Python 或 Filebeat,全部依赖都由容器提供。
唯一必须人工修改的配置是 Slack 凭据:编辑 pathway-src/alerts.py 顶部,把两个占位符替换成真实的频道 ID 与 Bot Token:
SLACK_ALERT_CHANNEL_ID = "你的频道ID" # 例如 C01XXXXXX
SLACK_ALERT_TOKEN = "你的Bot访问令牌" # 例如 xoxb-...
同时在 alerts.py 处理 license key:使用 Community 版本时,把 pw.set_license_key("demo-license-key-with-telemetry") 一行注释掉即可。
7.2 一键启动与验证
关联文档给出的启动流程为:
# 1. 在示例工程根目录启动全部容器(filebeat / zookeeper / kafka / pathway)
make
# 2. 连接到 Filebeat 容器
make connect
# 3. 在 Filebeat 容器内生成日志流
./generate_input_stream.sh
执行后,符合预期时 Slack 频道会立刻开始收到告警消息:脚本前 100 秒低频写入不会触发;进入高频写入阶段后,任何 1 秒窗口内日志数达到 5 条,Pathway 就会把告警状态置为 True 并通过回调推送 ACTIVE,状态回落时再推送 INACTIVE。
停止整个环境:
make stop
7.3 观察中间结果与调试技巧
关联文档提供了一个非常实用的调试手段:通过 pw.io.csv.write 把 log_table 落盘,在 Pathway 容器内直接查看原始日志流是否到达:
# 进入 Pathway 容器
make connect-pathway
# 查看输出的 CSV(先需在 alerts.py 中加入如下一行)
cat logs.csv
即在 alerts.py 的 pw.run() 之前加入:
pw.io.csv.write(log_table, "./logs.csv")
这样 pw.io.csv.write 会持续把 log_table 的增量写入 logs.csv(容器内),用来核对 Kafka 消费与时间解析是否正常;排查完成后删除该行并重建镜像即可。类似的思路也可用于观察 t_sliding_window 与 t_alert 的中间结果。
八、扩展:把链路改造成你需要的形态
这套模板在仓库中还有两处可直接对照的变体:
- 走 Logstash + Elasticsearch 的完整 ELK 版:若希望保留 Logstash 做过滤、并把统计结果写回 Elasticsearch 检索,可对照姊妹工程 examples/projects/realtime-log-monitoring/logstash-pathway-elastic(Filebeat → Logstash → Kafka → Pathway → Elasticsearch),其中 Pathway 侧使用
pw.io.elasticsearch.write(t_alert, index_name=...)输出,教程 7.realtime-log-monitoring.md 对两种场景都有完整配置。 - 告警策略自定义:本文判定逻辑只依赖
alert_threshold与sliding_window_duration两个常量。把二者改为通过环境变量注入、或把 5 条/秒的固定阈值替换为相对历史基线的动态阈值(例如与上一时段均值比较),即可沉淀为更通用的异常检测能力;Pathway 的增量更新特性保证了策略调整后下游结果会自动、无重复地刷新。
结语
通过这套 Filebeat + Kafka + Pathway + Slack 的端到端示例,可以看到流式日志监控的完整形态:既有 docker-compose.yml 级别的工程化编排,也有 alerts.py 中覆盖"读 Kafka → 解析时间 → 滑动窗口聚合 → 遗忘旧数据 → 订阅式输出告警"的完整 Pathway 实现。与传统按固定周期轮询的窗口方案相比,Pathway 的 windowby 是事件驱动的增量计算:窗口永远基于最新到达的数据维护,不重复计算、不遗漏数据点,配合 common_behavior 的 cutoff 遗忘机制还可让长跑任务保持有界内存。读者可直接克隆本示例仓库,填入自己的 Slack 凭据后 make 一键复现这条实时日志告警流水线。
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