能否在Elasticsearch Watch中实现Action间HTTP响应传递?
Elasticsearch Watcher动作链实现工单ID回写Elastic文档
问题背景
我需要通过Elasticsearch Watcher的Action在工单系统中创建工单,并将工单系统返回的工单ID写入对应的Elasticsearch告警文档。工单ID包含在创建工单的webhook响应里,但我不知道如何在后续Action中访问这个响应。原本考虑用Input Chaining,但因为需要遍历查询到的告警结果(这只能在Actions里用foreach实现),所以需要类似Action Chaining的机制,但不确定是否支持。当前的Watch代码如下:
{ "trigger": { "schedule": { "interval": "1m" } }, "input": { "search": { "request": { "search_type": "query_then_fetch", "indices": [ ".alerts-security.alerts-default" ], "rest_total_hits_as_int": true, "body": { "query": { "bool": { "filter": [ { "range": { "@timestamp": { "gte": "now-1m", "lte": "now" } } }, { "match": { "kibana.alert.workflow_status": "open" } } ] } } } } } }, "condition": { "compare": { "ctx.payload.hits.total": { "gt": 0 } } }, "actions": { "otobo_webhook": { "transform": { "script": { "source": """ ['items': ctx.payload.hits.hits.collect(alert -> [ 'id': alert._id, 'timestamp': alert._source['@timestamp'], 'rule_name': alert._source['kibana.alert.rule.name'], 'rule_note': alert._source['kibana.alert.rule.note'], 'host_name': alert._source['host.name'] ])] """, "lang": "painless" } }, "foreach": "ctx.payload.items", "max_iterations": 500, "webhook": { "scheme": "", "host": "", "port": , "method": "POST", "path": "", "params": {}, "headers": {}, "body": """ { "UserLogin": "", "Password": "", "Ticket": { "Title": "Elastic Alert (generated by rule {{ctx.payload.rule_name}})" }, "Article": { "CommunicationChannel": "Email", "From": "test@test.pl", "Subject": "Elastic Alert (generated by rule {{ctx.payload.rule_name}})", "Body": "<p>Rule \"{{ctx.payload.rule_name}}\" emitted an alert at {{ctx.payload.timestamp}} with a note: \"{{ctx.payload.rule_note}}\". Device: {{ctx.payload.host_name}}.</p>", "ContentType": "text/html charset=utf-8" } } """ } } } }
解决方案
Elasticsearch Watcher支持动作链(Action Chaining),可以通过action.chain类型串联多个动作步骤,并且能捕获每个步骤的响应供后续动作使用。针对你的场景,修改方案如下:
- 将原有的工单创建webhook动作作为动作链的第一个步骤,捕获其响应
- 添加第二个步骤,读取响应中的工单ID,更新对应的Elasticsearch告警文档
修改后的完整Watch配置:
{ "trigger": { "schedule": { "interval": "1m" } }, "input": { "search": { "request": { "search_type": "query_then_fetch", "indices": [ ".alerts-security.alerts-default" ], "rest_total_hits_as_int": true, "body": { "query": { "bool": { "filter": [ { "range": { "@timestamp": { "gte": "now-1m", "lte": "now" } } }, { "match": { "kibana.alert.workflow_status": "open" } } ] } } } } } }, "condition": { "compare": { "ctx.payload.hits.total": { "gt": 0 } } }, "actions": { "ticket_and_update_chain": { "action.chain": { "actions": [ // 步骤1:创建工单,保留原有的foreach和webhook配置 { "id": "create_ticket", "transform": { "script": { "source": """ ['items': ctx.payload.hits.hits.collect(alert -> [ 'id': alert._id, 'timestamp': alert._source['@timestamp'], 'rule_name': alert._source['kibana.alert.rule.name'], 'rule_note': alert._source['kibana.alert.rule.note'], 'host_name': alert._source['host.name'] ])] """, "lang": "painless" } }, "foreach": "ctx.payload.items", "max_iterations": 500, "webhook": { "scheme": "", "host": "", "port": , "method": "POST", "path": "", "params": {}, "headers": {}, "body": """ { "UserLogin": "", "Password": "", "Ticket": { "Title": "Elastic Alert (generated by rule {{ctx.payload.rule_name}})" }, "Article": { "CommunicationChannel": "Email", "From": "test@test.pl", "Subject": "Elastic Alert (generated by rule {{ctx.payload.rule_name}})", "Body": "<p>Rule \"{{ctx.payload.rule_name}}\" emitted an alert at {{ctx.payload.timestamp}} with a note: \"{{ctx.payload.rule_note}}\". Device: {{ctx.payload.host_name}}.</p>", "ContentType": "text/html charset=utf-8" } } """ } }, // 步骤2:更新Elasticsearch文档,写入工单ID { "id": "update_alert_with_ticket_id", "transform": { "script": { "source": """ // 关联原告警ID和工单ID:将create_ticket的响应与原items匹配 def updates = []; for (int i = 0; i < ctx.payload.items.size(); i++) { def alert = ctx.payload.items[i]; def response = ctx.actions.create_ticket.responses[i]; if (response.status == 200 && response.payload?.TicketID) { updates.add([ 'update': { '_index': '.alerts-security.alerts-default', '_id': alert.id, 'doc': { 'kibana.alert.workflow_status': 'ticket_created', 'ticket_id': response.payload.TicketID // 替换为工单系统返回的实际字段名 } } ]); } } return ['bulk_body': updates]; """, "lang": "painless" } }, "elasticsearch": { "api": "bulk", "body": "{{ctx.payload.bulk_body}}" } } ] } } } }
关键说明
action.chain会按顺序执行数组内的动作步骤,前一步的响应会被存储在ctx.actions.<步骤id>.responses中- 因为第一步使用了
foreach,responses是一个数组,每个元素对应一次迭代的webhook响应,顺序与ctx.payload.items一致 - 第二步通过Painless脚本遍历响应数组,提取工单ID,生成Elasticsearch Bulk更新请求,将工单ID写入原告警文档,并更新工单状态
内容的提问来源于stack exchange,提问作者Pietrek
相关产品推荐
相关产品推荐

