上周接了一个内部通知需求本来只是想把“查询天气”“判断是否要提醒”“发到 Slack”这三件事串起来结果发现大家嘴上说的 Workflow 其实差别很大。有人想用 Elasticsearch 8.15 之后自带的搜索工作流能力有人只是在找一个能定时跑脚本的编排器还有人只是想解决“怎么从 ES 查出结果再扔给 Slack”这种最朴素的场景。这篇文章我不会只贴代码而是把从需求拆解、环境准备、天气数据入库、查询判定、Slack 通知到最终封装成可调度工作流的完整过程讲清楚。项目标题里那几个词——Elasticsearch、Workflow、Slack——每个都有值得单独展开的坑我尽量用跑通过的实践来说。1. 先想清楚你要编排的到底是什么1.1 表面需求是“查天气”真实需求是“有条件的通知”直接拿一个脚本调天气 API然后把 temperature 拼成一段文本发到 Slack。这个过程五分钟能写完但它不算 Workflow因为它没有“决策”。实际工作中我们要的不是每天早上七点被一条“上海当前温度 12 度”的消息吵醒而是只在天气达到某个阈值时才收到提醒比如当前有中雨以上天气出门要带伞风速超过 30 km/h不适合户外骑行温度低于 5 度并且天气为雨夹雪需要提醒值班人员注意路况。所以这个需求真正困难的地方在于查询不是目的在什么条件下把查询结果投递给人才是目的。Workflow 的核心就是把“取数据、存数据、判定、通知”这四个步骤拆成独立环节同时保留一条清晰的执行链路。1.2 我当时为什么选择 Elasticsearch 做承载有人问过既然天气是从第三方 API 拿到的为什么还要先写进 Elasticsearch直接在内存里判断不行吗单次查询确实可以但如果接下来要做天气历史趋势分析、按城市维度横向对比、或者把天气数据和业务故障数据放在同一个看板里观察Elasticsearch 的价值就出来了。Elasticsearch 本身最强的是检索和聚合不是跑定时调度所以我的设计思路很朴素上游Open-Meteo 或者其他天气 API负责提供实时天气中游Elasticsearch负责保存时序化天气快照、提供查询和判断依据下游Slack Incoming Webhook负责把筛选后的消息推给对应频道。整条链路就是一个可复用的工作流。如果你想在 ES 原生能力里找到名为 Workflow 的独立模块8.15 以后的版本确实在逐步演进搜索工作流但目前更通用、更可控的方案还是用外部调度脚本把几个动作串起来。先把能稳定跑通的链路做出来再谈要不要迁移到原生 API。1.3 数据流和环节职责下面这张表是我在设计阶段给自己画的避免写到一半把职责搞混。阶段输入动作输出1. fetch城市经纬度调用天气 API结构化天气 JSON2. store结构化天气 JSON写入 ES weather_metrics 索引可检索的历史记录3. query城市名查询最新一条并套判定规则是否告警4. notify告警文本POST 到 Slack WebhookSlack 消息每一步都可以单独替换。比如今天用 Open-Meteo明天想换和风天气只需要改 fetch今天用 Slack明天想推到企业微信只需要改 notify。这就是把一个散装脚本整理成 Workflow 的最大收益。2. 准备一套能跑实验的 Elasticsearch 环境2.1 Docker Compose 是最省事的方案本机装过 ES 的人都懂依赖 JDK、内存配置、证书、权限哪个环节出问题都能卡半天。我这里直接使用 Docker Compose 起一个单节点集群集群版本选择 8.13.2。我知道很多人搜过 9.0.4 Windows 版本下载但实际项目里不必追新8.x 足够稳而且下面的配置思路在 9.x 上也通用。services: es: image: docker.elastic.co/elasticsearch/elasticsearch:8.13.2 container_name: es-weather-demo environment: - discovery.typesingle-node - xpack.security.enabledtrue - ELASTIC_PASSWORDyour_password - ES_JAVA_OPTS-Xms2g -Xmx2g ports: - 9200:9200 volumes: - es_data:/usr/share/elasticsearch/data volumes: es_data:启动命令很简单docker compose up -d docker compose logs -f es看到message: started一般就是起来了。这里有个新手容易忽略的问题8.x 默认开启安全认证并且 Elasticsearch 默认暴露的是 HTTPS 端口使用浏览器、curl 或者 Python 客户端时不能再用http://localhost:9200直连。2.2 在 Windows 上跑 Elasticsearch 的额外注意如果你没有 Docker 环境想直接用 Windows 压缩包跑也不是不行。解压后建议把目录放到不带中文和不带空格的路径下比如D:\elasticsearch-8.13.2否则启动时容易因为路径解析出问题。启动前至少要改两个地方。一是调整 JVM 堆内存。ES 默认堆内存可能偏小开发机建议打开config/jvm.options把-Xms和-Xmx改成相同的值比如 2g。这里强调相同值是为了避免运行中出现动态扩容导致停顿。二是确认本机端口没有占用。如果 9200 被占ES 会启动失败。Windows 下命令netstat -ano | findstr :9200没有输出就代表端口空闲。确认没问题后再执行bin\elasticsearch.bat不要双击这个 bat。因为一旦直接双击错误信息一闪而过你根本不知道是内存问题还是端口问题。在终端里启动至少能看到报错。等窗口稳定后另开一个 PowerShell 执行curl -k -u elastic:your_password https://localhost:9200因为本机是自签名证书-k用于跳过证书校验。如果看到tagline : You Know, for SearchES 就没问题了。2.3 用 Python 客户端完成连接验证后面代码会用到 Elasticsearch 官方的 Python 客户端先安装依赖pip install elasticsearch requests然后写一个最简单的健康检查from elasticsearch import Elasticsearch es Elasticsearch( https://localhost:9200, basic_auth(elastic, your_password), verify_certsFalse, request_timeout30 ) print(es.info())注意这里连接串写的是 HTTPS。如果只想在本地开发验证而不想处理证书可以使用verify_certsFalse但生产环境务必用 CA 证书或者 API Key不要整体关闭校验。3. 上游数据用 Open-Meteo 把天气拉下来写进 ES3.1 为什么不选那些看起来很复杂的天气服务天气 API 有很多选择但大家会发现很多服务要先注册、要领 API Key、还限制调用次数。Open-Meteo 是一个很适合做实验的免费天气接口不需要 Key直接通过 HTTP GET 就能拿到预报数据。它的 URL 大概长这样https://api.open-meteo.com/v1/forecast?latitude31.2304longitude121.4737current_weathertruelatitude 是纬度longitude 是经度current_weathertrue表示只返回当前天气而不是未来小时级预报。返回的 JSON 是这样{ latitude: 31.23, longitude: 121.47, current_weather: { temperature: 12.4, windspeed: 18.3, winddirection: 120, weathercode: 3, time: 2025-06-20T10:15 } }实际生产场景中要想清楚调用免费 API 是否有稳定性和数据授权问题但在技术 demo 和内部原型阶段它够用了。3.2 把天气快照写入 Elasticsearch 的正确姿势天气数据本质上是一种小体量时序数据。我建议不要裸调 API 后直接用结果发 Slack而是先写入索引这样后续想查看历史、想用 Kibana 画趋势图都有数据可用。先创建索引映射。天气字段里面存在不少数字如果不提前声明映射ES 会根据第一条数据自动推断。自动映射在 demo 里问题不大但我更建议手动创建映射因为后面查询和聚合会更可控。在 Kibana Dev Tools 中执行PUT /weather_metrics { mappings: { properties: { timestamp: { type: date }, city: { type: keyword }, temperature: { type: float }, windspeed: { type: float }, winddirection: { type: short }, weathercode: { type: keyword }, record_time: { type: keyword } } } }需要注意record_time 是天气接口返回的当地观测时间我把它设计成 keyword而不是 date。原因是 Open-Meteo 返回的时间字符串没有带时区如果强行映射为 dateES 会按 UTC 解析和国内本地时间直观对比时容易差 8 个小时。排序和筛选统一使用timestamp这个字段代表数据写入 Elasticsearch 的 UTC 时间语义清晰。Python 侧拉取并写入from datetime import datetime, timezone import requests from elasticsearch import Elasticsearch es Elasticsearch( https://localhost:9200, basic_auth(elastic, your_password), verify_certsFalse, request_timeout30 ) def fetch_weather(city_name: str, lat: float, lon: float) - dict: resp requests.get( https://api.open-meteo.com/v1/forecast, params{ latitude: lat, longitude: lon, current_weather: true }, timeout15 ) resp.raise_for_status() cur resp.json()[current_weather] return { timestamp: datetime.now(timezone.utc).isoformat(), city: city_name, temperature: cur[temperature], windspeed: cur[windspeed], winddirection: cur.get(winddirection), weathercode: str(cur[weathercode]), record_time: cur[time] } def store_weather(doc: dict) - str: doc_id f{doc[city]}-{doc[record_time]} return es.index( indexweather_metrics, iddoc_id, documentdoc )[_id]这里值得展开一下 doc_id 的设计。如果我们每天定时跑任务Open-Meteo 在某个观测时刻返回的天气数据是同一个天气事件。使用城市-观测时间作为 document id可以在同一时刻重复跑任务时覆盖同一条记录而不是反复新增。批量采集多个城市时也不需要额外装什么 bulk 插件官方客户端自带批量处理能力数据量大再改成 elasticsearch.helpers.bulk 就行。3.3 先落库存查询到底多了哪些好处很多人觉得多了一步写 ES 很麻烦我实际用下来觉得有三个价值。第一历史数据可追溯。如果某天 Slack 推了一条“暴雨预警”事后复盘时可以直接在 ES 里查当时的天气数据确认告警是否合理。第二判断逻辑和采集逻辑解耦。天气采集可能需要每十分钟一次但通知只需要每天一次或者达到阈值才触发。如果每次采集后都直接执行通知逻辑后面做降噪会非常痛苦。把数据先存下来通知环节按需查询代码会清爽很多。第三可以和已有业务数据联动。比如把天气数据和当天的订单量、故障工单量放到同一张可视化看板观察天气对业务指标的影响。这在纯 API 脚本方案里很难做但数据一旦进入 ES就能用现成的聚合能力去分析。4. 让查询决定要不要发ES 查询和天气告警逻辑4.1 怎么查出某城市最新一条天气记录ES 查询要返回“最新一条”不能直接 size1 然后碰运气必须按时间排序。查询语句如下GET /weather_metrics/_search { size: 1, query: { bool: { filter: [ { term: { city: Shanghai } } ] } }, sort: [ { timestamp: { order: desc } } ] }Python 侧封装成函数后就能在后续 Workflow 中复用了。def get_latest_weather(city: str) - dict | None: resp es.search( indexweather_metrics, size1, query{term: {city: city}}, sort[{timestamp: {order: desc}}], source[city, temperature, windspeed, weathercode, record_time] ) hits resp[hits][hits] if not hits: return None return hits[0][_source]这里有个小优化通过source只请求需要的字段避免把整个文档完整拉回来。数据量小时不明显但数据量上来了之后能减少网络传输和内存压力。4.2 把天气编码翻译成“要不要提醒”Open-Meteo 的 weathercode 是一组数字编码想要理解这些编码我对着文档归纳了一张简表。这张表不需要特别精确因为最终判定规则可以根据业务调整。weathercode含义0晴1-2少云到多云3阴45/48雾51-57毛毛雨61-67雨71-77雪80-82阵雨95-99雷暴我当时的告警规则是雨相关编码命中 61、63、65、80、81、82 时提醒带伞命中 95 到 99 时因为可能伴随雷电优先级提高另外单独把风速大于等于 30 km/h 当作另一个独立条件。这个 30 不是公式算出来的是结合户外活动经验定的不同团队完全应该按自己的场景改。Python 判定函数可以这样写RAINY_CODES {61, 63, 65, 80, 81, 82} def should_alert(source: dict) - bool: code source.get(weathercode, ) wind float(source.get(windspeed, 0) or 0) return code in RAINY_CODES or wind 304.3 组装一段能看懂的 Slack 文本判定通过后不能直接把 JSON 原封不动发到 Slack那样群里全是代码噪音。最好整理成人类可读的文本。我当时的做法是做一层简单的编码翻译再做一次字符串拼接。CODE_LABEL { 0: 晴天, 1: 少云, 3: 阴天, 61: 小雨, 63: 中雨, 65: 大雨, 80: 阵雨, 95: 雷阵雨 } def build_message(source: dict) - str: city source[city] code source.get(weathercode, unknown) label CODE_LABEL.get(code, f未知天气码({code})) temp source[temperature] wind source[windspeed] record_time source.get(record_time, ) reason [] if code in RAINY_CODES: reason.append(有降雨) if float(wind) 30: reason.append(风速偏大) return ( f{city} 天气提醒\n f观测时间{record_time}\n f天气{label}\n f温度{temp}℃\n f风速{wind} km/h\n f提醒原因{.join(reason)} )我在实际调试中发现一个坑Open-Meteo 返回的 weathercode 是数字比如 61。如果索引映射里定义为 keywordPython 侧也方便比较但如果你从 JSON 里直接取数后发现是 int而你的判定集合里写的是字符串两者永远不相等。建议在 fetch 阶段就做str(cur[weathercode])统一转字符串。数据来源多时类型不统一是告警逻辑失效的重灾区。5. 发 Slack 通知连接 Webhook 这一步没想象中简单5.1 创建 Slack Incoming WebhookSlack 发送消息最快速的方式是 Incoming Webhook。打开 Slack 的管理后台进入 API 页面创建一个 App然后在 Incoming Webhooks 开关下添加一个 Webhook选好目标频道Slack 会给你生成一个形如下面的长地址https://hooks.slack.com/services/T00000000/B00000000/XXXXXXXXXXXXXXXXXXXXXXXX这里的 Webhook 相当于一个带写入权限的 URL谁拿到这个地址谁就能往对应频道发消息。因此不要把 Webhook 直接写死在仓库里。我会建议用一个.env文件保存然后在 Python 中读取SLACK_WEBHOOK_URLhttps://hooks.slack.com/services/xxxx5.2 用 Python 发送一条消息Slack 的 Incoming Webhook 要求 POST 一个 JSON核心字段就是text。最简单的代码import os def send_slack(text: str) - bool: webhook_url os.environ.get(SLACK_WEBHOOK_URL) if not webhook_url: print(SLACK_WEBHOOK_URL 未配置发送失败) return False resp requests.post( webhook_url, json{text: text}, timeout15 ) if resp.status_code ! 200: print(fSlack 返回异常{resp.status_code} {resp.text}) return False return True正常返回是 HTTP 200 并且响应体为ok。很多教程只贴一个 curl 命令但发消息之前需要留意两点。第一Slack 消息内容里的换行符必须是字符串里的真实\n如果从某个配置系统里读多行文本注意别把\\n原样传过去。第二如果消息内容后续要放入 JSON尽量使用json{text: text}而不是手动拼接 JSON 字符串。手动拼接时一旦文本里出现引号整个 payload 就会格式错误。5.3 消息发送的“成功”要如何定义我之前犯过的一个错误是只看 HTTP 200 就认为发送成功。Incoming Webhook 基本可以这样做但严格一点还要确认响应体里没有 error。响应体正常是一个单词ok如果出现类似invalid_payload的报错说明提交的 JSON 结构不对。如果遇到 HTTP 429说明触发了 Slack 的限流需要退避重试。if resp.status_code 429: print(Slack 限流建议稍后重试)Slack 本身是外部服务网络超时和限流都可能发生。在 Workflow 中加入简单的超时和重试逻辑能有效减少“脚本没报错但消息没发出”的悬案。6. 把整个流程封装成一个可调度的 Workflow6.1 先把各种函数整理成一个类到这一步我们已经有了 fetch_weather、store_weather、get_latest_weather、should_alert、build_message、send_slack 这些函数。如果继续在脚本里从上到下顺序调用并非不能用但后续要加城市循环、加条件判断会越来越乱。我建议封装成一个 WeatherWorkflow 类把四个核心步骤暴露成有名字的方法class WeatherWorkflow: def __init__(self, es_client, slack_webhook_url): self.es es_client self.slack_webhook_url slack_webhook_url def run(self, city: str, lat: float, lon: float) - None: doc self._fetch_weather(city, lat, lon) self._store_weather(doc) latest self._get_latest_weather(city) if not latest: print(没有查到最新天气记录) return if not should_alert(latest): print(f{city} 当前天气不满足提醒条件跳过 Slack 发送) return text build_message(latest) self._send_slack(text) def _fetch_weather(self, city, lat, lon): return fetch_weather(city, lat, lon)这里的优点不是代码变少而是当执行顺序需要变化时你只要改 run 方法的行序。有人可能觉得这个类叫 Workflow 有点“重”但把整条业务链路的入口收敛成一个 run() 方法远比在多个脚本之间复制粘贴更便于维护。6.2 调度方式每天定时执行让 workflow 跑起来需要一个触发机制。Linux 下用 cron 就可以写一个工作目录下的执行脚本0 7 * * * cd /opt/weather_job python -m workflow.main如果想在 Windows 上做定时可以用任务计划程序。我简单说下命令行方式这也是很多人搜“windows 启动 elasticsearch”之后会遇到的任务配置问题schtasks /create /tn weather_slack /tr python C:\weather_job\main.py /sc daily /st 07:00不过在 Windows 上直接写python有可能找不到解释器因为系统没有把 Python 路径加入 PATH。先执行where python查一下完整路径再放进任务计划里比较稳妥。这一类问题看起来小实际排查起来很费时间。6.3 避免重复告警加一个状态标记定时任务跑起来以后一定会遇到重复告警问题。比如当天 7 点调度了一次天气正好是中雨发送了一条 Slack如果 8 点又调度了一次天气还是中雨理论上这条消息不需要再次发送否则群里会变得很吵。我推荐的方案是在 ES 里增加一个 weather_alert_state 索引记录每个城市已发送过的天气事件标识。发送前先查询是否已经发过。标识可以用城市-天气观测时间-天气编码生成。def _already_notified(city: str, record_time: str, code: str) - bool: marker f{city}-{record_time}-{code} return es.exists(indexweather_alert_state, idmarker) def _mark_notified(city: str, record_time: str, code: str) - None: marker f{city}-{record_time}-{code} es.index( indexweather_alert_state, idmarker, document{marker: marker, notified_at: datetime.now(timezone.utc).isoformat()} )发送成功后调用_mark_notified发送前调用_already_notified。这样相同天气事件重复触发时不会每条都推到 Slack。这里有一个取舍如果先把“未通知”标记为“已通知”再去调 Slack万一 Slack 发送失败这条提醒就永久丢失了所以我实际是发送成功后再写状态标记极端情况下可能因为网络超时造成重复发送一次但比漏掉重要提醒更容易接受。6.4 单独聊一句ES 原生 Workflow API 的定位项目标题里出现 Workflow 时有人期待的是 ES 原生 Workflow API。在 8.15 之后的版本里Elasticsearch 确实在搜索侧引入了工作流能力用来把查询、推理、后处理这些动作编排到一个请求里典型场景是搜索增强、向量召回和 LLM 结合。但它解决的是“一次搜索请求内多个 Elasticsearch 动作的编排”不是