You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否在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类型串联多个动作步骤,并且能捕获每个步骤的响应供后续动作使用。针对你的场景,修改方案如下:

  1. 将原有的工单创建webhook动作作为动作链的第一个步骤,捕获其响应
  2. 添加第二个步骤,读取响应中的工单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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 07:24:59