首页
/ Pathway 实时日志监控实战:Filebeat + Kafka + Pathway 滑动窗口告警并推送 Slack

Pathway 实时日志监控实战:Filebeat + Kafka + Pathway 滑动窗口告警并推送 Slack

2026-09-07 10:36:47作者:舒璇辛Bertina

实时监控服务器日志(如 Nginx 访问日志)并第一时间发现流量异常,是流式处理最典型的落地场景之一。本文基于开源仓库 pathwayexamples/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 中):

  1. 字段抽取:从 Filebeat 生成的 JSON 消息中取出 @timestampmessage 两个字段;
  2. 时间解析:把 ISO 8601 格式的字符串时间转换为真正的(可运算的)时间戳;
  3. 滑动窗口过滤:只保留最近 X 秒(默认 X=1s)内的日志,其中"当前时间"以窗口内最后一条日志的时间戳为参照;
  4. 阈值判定:若窗口内日志数超过 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 消息包含很多元数据字段,本示例只关心 @timestampmessage 两列,其余字段可忽略。由于列名 @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_sslSCRAM-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=10msduration=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_additionTrue 表示该行进入表(告警状态变为 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/Dockerfilepython:3.10 基础上安装 pathwayrequestspython-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]

值得注意的三点实现细节:

  1. topic 自动创建:Kafka 容器的 command 在启动后延迟 15 秒,用 kafka-topics --create 创建名为 logs 的单分区主题,再拉起 Kafka 主进程。若想换主题名,需要同步修改这一行、alerts.py 里的 topic="logs" 以及 Filebeat 的 output.kafka 配置。
  2. 容器内互联用服务名:Pathway 中 rdkafka_settingsbootstrap.servers 必须写 kafka:9092(compose 服务名),Filebeat 的 Kafka output host 同样写 kafka:9092,而不是 localhost
  3. 服务间依赖kafka 依赖 zookeeperpathway 依赖 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 stopdocker 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.writelog_table 落盘,在 Pathway 容器内直接查看原始日志流是否到达:

# 进入 Pathway 容器
make connect-pathway
# 查看输出的 CSV(先需在 alerts.py 中加入如下一行)
cat logs.csv

即在 alerts.pypw.run() 之前加入:

pw.io.csv.write(log_table, "./logs.csv")

这样 pw.io.csv.write 会持续把 log_table 的增量写入 logs.csv(容器内),用来核对 Kafka 消费与时间解析是否正常;排查完成后删除该行并重建镜像即可。类似的思路也可用于观察 t_sliding_windowt_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_thresholdsliding_window_duration 两个常量。把二者改为通过环境变量注入、或把 5 条/秒的固定阈值替换为相对历史基线的动态阈值(例如与上一时段均值比较),即可沉淀为更通用的异常检测能力;Pathway 的增量更新特性保证了策略调整后下游结果会自动、无重复地刷新。

结语

通过这套 Filebeat + Kafka + Pathway + Slack 的端到端示例,可以看到流式日志监控的完整形态:既有 docker-compose.yml 级别的工程化编排,也有 alerts.py 中覆盖"读 Kafka → 解析时间 → 滑动窗口聚合 → 遗忘旧数据 → 订阅式输出告警"的完整 Pathway 实现。与传统按固定周期轮询的窗口方案相比,Pathway 的 windowby 是事件驱动的增量计算:窗口永远基于最新到达的数据维护,不重复计算、不遗漏数据点,配合 common_behavior 的 cutoff 遗忘机制还可让长跑任务保持有界内存。读者可直接克隆本示例仓库,填入自己的 Slack 凭据后 make 一键复现这条实时日志告警流水线。

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

项目优选

收起
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