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

Logstash消费Azure EventHub时遇Binary转换缺失错误求助

Logstash从Azure EventHub拉取数据时MissingConverterException问题解决

问题描述

配置Logstash从Azure EventHub拉取事件时,日志出现org.logstash.MissingConverterException错误,核心报错为Missing Converter handling for full class name=org.apache.qpid.proton.amqp.Binary。EventHub中存储的是JSON数据,发布消息使用Spring Kafka的JsonSerializer,配置如下:

spring.kafka.producer:
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

错误日志片段:

[2025-09-10T11:26:34,355][WARN ][com.microsoft.azure.eventprocessorhost.PartitionPump][main][b71544191bfb680f60e090ab62fea980d20ad24b8fff4ded3dba1fe383a8cd0d] host logstash-a36d07f7-f854-4835-835e-e9137f7484ee: 0: Got exception from onEvents
org.logstash.MissingConverterException: Missing Converter handling for full class name=org.apache.qpid.proton.amqp.Binary, simple name=Binary
    at org.logstash.Valuefier.fallbackConvert(Valuefier.java:122) ~[logstash-core.jar:?]
    at org.logstash.Valuefier.convert(Valuefier.java:100) ~[logstash-core.jar:?]
    at org.logstash.ConvertedMap.newFromMap(ConvertedMap.java:111) ~[logstash-core.jar:?]
    at org.logstash.Valuefier.lambda$initConverters$15(Valuefier.java:182) ~[logstash-core.jar:?]
    at org.logstash.Valuefier.convert(Valuefier.java:98) ~[logstash-core.jar:?]
    at org.logstash.ext.JrubyEventExtLibrary$RubyEvent.safeValueifierConvert(JrubyEventExtLibrary.java:350) ~[logstash-core.jar:?]
    at org.logstash.ext.JrubyEventExtLibrary$RubyEvent.ruby_set_field(JrubyEventExtLibrary.java:122) ~[logstash-core.jar:?]

原因分析

Spring Kafka的JsonSerializer默认会将Java对象序列化为二进制格式(而非纯JSON字符串),导致Logstash的Azure EventHub输入插件接收到Binary类型数据,而Logstash默认没有对应转换器处理该类型,因此抛出异常。

解决方法

方案1:改用StringSerializer手动序列化JSON

修改Spring Kafka生产者配置,使用StringSerializer,并在业务代码中将对象手动序列化为JSON字符串:

  1. 修改配置文件:
spring.kafka.producer:
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
  1. 业务代码示例(Java):
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.kafka.core.KafkaTemplate;

// 注入KafkaTemplate和ObjectMapper
private final KafkaTemplate<String, String> kafkaTemplate;
private final ObjectMapper objectMapper;

// 发送消息时手动序列化
public void sendJsonMessage(String topic, Object data) throws Exception {
    String jsonStr = objectMapper.writeValueAsString(data);
    kafkaTemplate.send(topic, jsonStr);
}

方案2:调整JsonSerializer配置输出纯JSON字符串

若需保留JsonSerializer,可通过配置禁用类型信息输出,确保生成纯JSON字符串:

spring.kafka.producer:
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      properties:
        # 禁用类型信息写入
        spring.json.serialize.to.json.with.type: false
        # 指定默认序列化类型为字符串
        spring.json.serialize.default.type: java.lang.String

方案3:确保Logstash正确解析JSON

在Logstash的配置文件中添加json过滤器,将接收到的字符串解析为JSON结构:

filter {
  json {
    source => "message"
    # 可选:若解析失败保留原始消息
    tag_on_failure => ["json_parse_failure"]
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 09:13:10