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

Logstash解析XML关联JDBC数据存入ES时数量字段未生效排查

问题描述

使用Logstash 8.7与Elasticsearch 8.7,需将SQLite数据库products表(含ID、name、description、price字段)的数据存入ES索引,同时解析stocks.xml文件补充缺失的quantity字段(productid对应数据库ID),但目前quantity字段未成功存入索引,尝试两种配置均无效。

第一次尝试的logstash.conf

# Sample Logstash configuration for creating a simple
# Beats -> Logstash -> Elasticsearch pipeline.

input {
  beats {
    port => 5044
  }

  tcp {
    port => 50000
  }

  jdbc {
      jdbc_driver_library => "sqlite-jdbc-3.6.7.jar"
      jdbc_driver_class => "org.sqlite.JDBC"
      jdbc_connection_string => "jdbc:sqlite:products.db"
      jdbc_user => ""
      jdbc_password => ""
      statement => "SELECT * FROM products"
  }

  file {
      path => "stocks.xml"
      start_position => "beginning"
      sincedb_path => "/dev/null"
      codec => multiline {
        pattern => "^<product>"
        negate => true
        what => "previous"
      }
  }
}

filter {
  xml {
    source => "message"
    target => "parsed_xml"
    store_xml => false
    xpath => [
      "//product/productid/text()", "product_id",
      "//product/quantity/text()", "quantity"
    ]
  }
}

output {
  elasticsearch {
    hosts => ["http://localhost:9200"]
    index => "products"
    document_id => "%{id}"
  }
  stdout {
    codec => rubydebug
  }
}

stocks.xml文件

<?xml version="1.0" encoding="UTF-8"?>
<stock>
    <product>
        <productid>1</productid>
        <quantity>87</quantity>
    </product>
    <product>
        <productid>2</productid>
        <quantity>54</quantity>
    </product>
    <product>
        <productid>3</productid>
        <quantity>24</quantity>
    </product>
    <product>
        <productid>4</productid>
        <quantity>13</quantity>
    </product>
    <product>
        <productid>5</productid>
        <quantity>22</quantity>
    </product>
</stock>

第二次尝试的logstash.conf

input {
  beats {
    port => 5044
  }

  tcp {
    port => 50000
  }

  jdbc {
    jdbc_driver_library => "sqlite-jdbc-3.6.7.jar"
    jdbc_driver_class => "org.sqlite.JDBC"
    jdbc_connection_string => "jdbc:sqlite:products.db"
    jdbc_user => ""
    jdbc_password => ""
    statement => "SELECT * FROM products"
  }

  file {
    path => "stocks.xml"
    start_position => "beginning"
    sincedb_path => "/dev/null"
    codec => multiline {
      pattern => "<product>"
      negate => true
      what => "previous"
    }
  }
}

filter {
  xml {
    source => "message"
    target => "parsed_xml"
    store_xml => false
    force_array => false
    xpath => [
      "//product/productid/text()", "product_id",
      "//product/quantity/text()", "quantity"
    ]
  }
}

output {
  elasticsearch {
    hosts => ["http://localhost:9200"]
    index => "products"
    document_id => "%{id}"
  }
  stdout {
    codec => rubydebug
  }
}
解决思路

核心问题是JDBC输入和XML输入是独立数据流,未将两者关联,导致SQLite数据与XML的quantity字段未合并到同一条文档中。以下是具体修复方案:

1. 修正XML输入的多行解析逻辑

当前multiline配置无法正确拆分单个<product>节点,修改file输入的codec,确保每个<product>块被合并为单条事件:

file {
  path => "stocks.xml"
  start_position => "beginning"
  sincedb_path => "/dev/null"
  codec => multiline {
    pattern => "^</product>"
    negate => false
    what => "previous"
  }
}

2. 方案一:预生成字典+translate过滤器关联数据

2.1 导出XML数据为JSON字典

先运行临时配置生成stock_quantity.json:

input {
  file {
    path => "stocks.xml"
    start_position => "beginning"
    sincedb_path => "/dev/null"
    codec => multiline {
      pattern => "^</product>"
      negate => false
      what => "previous"
    }
  }
}

filter {
  xml {
    source => "message"
    xpath => [
      "//product/productid/text()", "product_id",
      "//product/quantity/text()", "quantity"
    ]
    remove_field => ["message", "@timestamp", "@version"]
  }
  mutate {
    convert => { "product_id" => "integer" }
    convert => { "quantity" => "integer" }
  }
}

output {
  file {
    path => "stock_quantity.json"
    codec => line { format => '{"%{product_id}": "%{quantity}"}' }
  }
}

生成的字典文件内容示例:

{"1": "87"}
{"2": "54"}
{"3": "24"}

2.2 主配置中关联数据

修改主配置,移除XML输入,添加translate过滤器匹配id字段补充quantity:

input {
  beats {
    port => 5044
  }

  tcp {
    port => 50000
  }

  jdbc {
    jdbc_driver_library => "sqlite-jdbc-3.6.7.jar"
    jdbc_driver_class => "org.sqlite.JDBC"
    jdbc_connection_string => "jdbc:sqlite:products.db"
    jdbc_user => ""
    jdbc_password => ""
    statement => "SELECT * FROM products"
  }
}

filter {
  translate {
    source => "id"
    target => "quantity"
    dictionary_path => "stock_quantity.json"
    fallback => 0 # 无匹配时设置默认值
  }
  mutate {
    convert => { "quantity" => "integer" } # 转换为数字类型
  }
}

output {
  elasticsearch {
    hosts => ["http://localhost:9200"]
    index => "products"
    document_id => "%{id}"
  }
  stdout {
    codec => rubydebug
  }
}

3. 方案二:使用aggregate过滤器实时关联

无需预生成字典,在Logstash运行时构建内存映射:

input {
  jdbc {
    jdbc_driver_library => "sqlite-jdbc-3.6.7.jar"
    jdbc_driver_class => "org.sqlite.JDBC"
    jdbc_connection_string => "jdbc:sqlite:products.db"
    jdbc_user => ""
    jdbc_password => ""
    statement => "SELECT * FROM products"
    add_field => { "type" => "product" }
  }

  file {
    path => "stocks.xml"
    start_position => "beginning"
    sincedb_path => "/dev/null"
    codec => multiline {
      pattern => "^</product>"
      negate => false
      what => "previous"
    }
    add_field => { "type" => "stock" }
  }
}

filter {
  if [type] == "stock" {
    xml {
      source => "message"
      xpath => [
        "//product/productid/text()", "product_id",
        "//product/quantity/text()", "quantity"
      ]
    }
    mutate {
      convert => { "product_id" => "integer" }
      convert => { "quantity" => "integer" }
      remove_field => ["message", "@timestamp", "@version"]
    }
    # 将库存数据存入内存,key为product_id
    aggregate {
      task_id => "%{product_id}"
      code => "map['quantity'] = event.get('quantity'); event.cancel()"
    }
  }

  if [type] == "product" {
    # 关联内存中的库存数据
    aggregate {
      task_id => "%{id}"
      code => "event.set('quantity', map['quantity'] || 0)"
      timeout => 30 # 等待30秒确保库存数据加载完成
    }
  }
}

output {
  elasticsearch {
    hosts => ["http://localhost:9200"]
    index => "products"
    document_id => "%{id}"
  }
  stdout {
    codec => rubydebug
  }
}

4. 验证要点

  • 查看stdout的rubydebug输出,确认JDBC事件是否包含quantity字段
  • 检查XML解析后的事件是否正确提取product_id和quantity
  • 确保id(SQLite)与product_id(XML)数据类型一致(均为整数),否则匹配失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 10:06:56