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

如何在MarkLogic中用反向查询与Kafka交互及触发告警?

MarkLogic反向查询结合Kafka消息处理方案

我来帮你梳理下这两个问题的具体实现方案,结合你已经在用的xdmp:http-post方式,直接上干货:


一、用反向查询向Kafka发送实际消息

反向查询(cts:reverse-query)的核心是找到引用了特定文档/满足特定条件的关联文档,再把这些文档的内容发送到Kafka。具体步骤如下:

  1. 构造反向查询条件:根据你的业务需求,定义要匹配的反向查询规则(比如找引用了某个URI文档的所有关联文档)
  2. 提取目标数据:执行查询后,从匹配的文档中提取需要发送的内容
  3. 格式化Kafka请求:按照Kafka REST Proxy的要求组装JSON payload,通过xdmp:http-post发送

示例代码

xquery version "1.0-ml";

(: 1. 构造反向查询:查找所有引用了"/docs/important-source.xml"的文档 :)
let $reverse-query := cts:reverse-query(cts:document-query("/docs/important-source.xml"))

(: 2. 执行查询并提取需要发送的内容 :)
let $matching-docs := cts:search(fn:doc(), $reverse-query)
let $kafka-records := 
  for $doc in $matching-docs
  return {
    "value": {
      "doc-uri": xdmp:node-uri($doc),
      "full-content": fn:string($doc)
    }
  }

(: 3. 格式化为Kafka要求的结构并发送 :)
let $kafka-payload := xdmp:to-json({
  "records": $kafka-records
})
return xdmp:http-post(
  "http://localhost:8082/topics/MyTopicName",
  <options xmlns="xdmp:http">
    <data>{$kafka-payload}</data>
    <headers>
      <content-type>application/vnd.kafka.json.v1+json</content-type>
    </headers>
  </options>
)

二、捕获传入消息,匹配反向查询后发送告警到Kafka

你已经实现了基础告警功能,现在只需要扩展逻辑:捕获传入消息→执行反向查询→判断匹配→发送原始消息全文。这里分两种场景处理:

场景1:监听MarkLogic文档写入(用触发器)

如果你的“传入消息”是指写入MarkLogic的文档,可以用后置触发器自动触发逻辑:

触发器关联的XQuery模块代码

xquery version "1.0-ml";

(: 触发器传入的变量:刚写入/更新的文档URI :)
declare variable $trigger-uri as xs:string external;
let $incoming-doc := fn:doc($trigger-uri)

(: 定义反向查询规则:比如查找引用了当前文档的告警关联文档 :)
let $alert-reverse-query := cts:reverse-query(cts:document-query($trigger-uri))
let $matched-docs := cts:search(fn:doc(), $alert-reverse-query)

(: 若有匹配结果,发送原始文档全文到告警主题 :)
return
  if (fn:exists($matched-docs)) then
    xdmp:http-post(
      "http://localhost:8082/topics/AlertTopic",
      <options xmlns="xdmp:http">
        <data>{
          xdmp:to-json({
            "records": [{
              "value": {
                "alert-type": "Reverse Query Match",
                "original-uri": $trigger-uri,
                "original-content": fn:string($incoming-doc),
                "matched-count": fn:count($matched-docs)
              }
            }]
          })
        }</data>
        <headers>
          <content-type>application/vnd.kafka.json.v1+json</content-type>
        </headers>
      </options>
    )
  else
    ()

场景2:接收外部传入的消息(用REST端点)

如果消息是通过REST API传入的,可以写一个自定义REST端点处理:

REST端点示例代码

xquery version "1.0-ml";
declare variable $input as node() external;

(: 从传入消息中提取关键标识,构造反向查询 :)
let $target-id := $input//reference-id/text()
let $reverse-query := cts:reverse-query(cts:element-value-query(xs:QName("reference-id"), $target-id))
let $matches := cts:search(fn:doc(), $reverse-query)

return
  if (fn:exists($matches)) then
    xdmp:http-post(
      "http://localhost:8082/topics/AlertTopic",
      <options xmlns="xdmp:http">
        <data>{
          xdmp:to-json({
            "records": [{
              "value": {
                "alert-message": "触发反向查询匹配告警",
                "original-message": fn:string($input)
              }
            }]
          })
        }</data>
        <headers>
          <content-type>application/vnd.kafka.json.v1+json</content-type>
        </headers>
      </options>
    )
  else
    <result>无匹配结果,未发送告警</result>

额外注意事项

  • 确保Kafka REST Proxy(kafka-rest)已正常启动,端口与代码中的8082一致
  • 若发送大文档,需调整Kafka的message.max.bytes配置,避免消息被截断
  • 可以添加try/catch块处理HTTP请求失败的情况,比如记录错误日志或重试

内容的提问来源于stack exchange,提问作者Dee

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:23:09