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
相关产品推荐
相关产品推荐

