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

使用Logstash向Kafka发送AVRO消息时遇SSLError证书验证失败如何解决?

问题描述

我有一个仅能消费AVRO编码消息的Kafka消费者,目前正尝试实现生产者,将系统日志文件中的日志以AVRO消息形式发送到Kafka主题,计划通过Logstash和Filebeat完成读取与发送。

此前配置JSON格式的Kafka输出时可正常工作,能从主题消费JSON消息,配置如下:

output {
  kafka {
    topic_id => "topic"
    bootstrap_servers => "servers"
    client_id => "id"
    security_protocol => "SASL_SSL"
    sasl_mechanism => "SCRAM-SHA-512"
    sasl_jaas_config => "org.apache.kafka.common.security.scram.ScramLoginModule required username='username' password='password';"
    ssl_truststore_location => "/etc/kafka.client.truststore.jks"
    ssl_truststore_password => "password"
    value_serializer => "org.apache.kafka.common.serialization.ByteArraySerializer"
    codec => json
  }
}

但实际需发送AVRO消息,改用包含avro_schema_registry的输出配置后,Logstash报错停止:

[logstash.javapipeline][main] Pipeline worker error, the pipeline will be stopped {:pipeline_id=>"main", :error=>"(SSLError) certificate verify failed", :exception=>Java::OrgJrubyExceptions::StandardError...

新配置如下:

output {
  kafka {
    topic_id => "topic"
    bootstrap_servers => "servers"
    client_id => "id"
    security_protocol => "SASL_SSL"
    sasl_mechanism => "SCRAM-SHA-512"
    sasl_jaas_config => "org.apache.kafka.common.security.scram.ScramLoginModule required username='username_AAA' password='password_AAA';"
    ssl_truststore_location => "/etc/kafka.client.truststore.jks"
    ssl_truststore_password => "password"
    value_serializer => "org.apache.kafka.common.serialization.ByteArraySerializer"
    codec => avro_schema_registry {
      endpoint => "endpoint"
      schema_id => 1111
      register_schema => false
      username => "username_AAA"
      password => "password_AAA"
    }
  }
}

我的消费者依赖Schema Registry,要求消息包含编码的schema id,因此生产者需传递该id。但为何已配置SSL信任库仍出现证书验证失败错误?为何需设置Schema Registry端点?能否直接传递含编码schema id的消息与字符串格式的schema?是否有无需avro_schema_registry的替代方案?


问题解答

1. 证书验证失败的原因

你当前配置的SSL信任库仅针对Kafka集群的连接生效,但avro_schema_registry codec需要单独连接Schema Registry服务,这个连接的SSL信任配置并未在现有配置中指定。也就是说,Logstash连接Schema Registry时使用的是默认信任根,而非你提供的kafka.client.truststore.jks,导致证书验证失败。

需要给avro_schema_registry codec补充SSL相关配置:

codec => avro_schema_registry {
  # 原有配置...
  ssl_truststore_location => "/etc/kafka.client.truststore.jks"
  ssl_truststore_password => "password"
}

如果Schema Registry的证书不在你的信任库中,需先将证书导入到kafka.client.truststore.jks,确保证书链完整。

2. 为何需要设置Schema Registry端点

即使指定了schema_id,avro_schema_registry codec仍需连接Schema Registry服务拉取对应的schema定义——它需要用这个schema来将Logstash事件数据编码成AVRO格式。没有端点的话,无法获取schema的具体结构,也就无法完成编码。

3. 能否直接传递schema id和字符串schema?

不行。AVRO消息有固定的二进制结构:前1字节是魔法值,接下来4字节是schema ID,最后是AVRO编码的消息体。直接拼接schema id和字符串schema不符合该规范,消费者无法解析。

4. 无需avro_schema_registry的替代方案

  • 自定义Ruby脚本编码:在Logstash的filter阶段使用ruby插件,提前加载本地AVRO schema文件,手动将事件数据编码为符合规范的AVRO二进制消息(包含魔法值、schema ID和编码内容),然后将codec设为plain、value_serializer用ByteArraySerializer直接输出到Kafka。
  • Filebeat直接发送AVRO:如果日志格式可匹配AVRO schema,可配置Filebeat的avro编码器搭配kafka输出,直接完成编码和发送,无需Logstash中转。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 08:27:40