使用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

