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

高吞吐场景下Java Kafka消费者字段合并为数组反序列化失败

Kafka消费者反序列化异常问题分析与解决

问题背景

使用Java开发Kafka消费者,期望接收单字段结构的JSON事件,正常格式示例:

{"cs-uri": "tcp://teams.microsoft.com:443/","cs(Referer)": "-","cs-auth-group": "/","cs-uri-extension": "-","cs-bytes": "13941","username": "-","src": "10.3.41.215", "dstport": "443","event_time": "Dec 1 06:46:09",...}

对应的实体类ProxyEvent定义如下:

public class ProxyEvent {

    @JsonProperty("cs-uri-scheme")
    String CsUriScheme;
    @JsonProperty("devicetime")
    String DeviceTime;
    @JsonProperty("dst")
    String Dst;
    @JsonProperty("event_source_id")
    String EventSourceId;
    @JsonProperty("log_type")
    String LogType;
    @JsonProperty("time-taken")
    String TimeTaken;
    @JsonProperty("rs(Content-Type)")
    String RsContentType;
    @JsonProperty("dstport")
    String Dstport;
    @JsonProperty("sc-status")
    String ScStatus;
    @JsonProperty("cs-uri-extension")
    String CsUriExtension;
    @JsonProperty("srcport")
    String Srcport;
}

当数据吞吐超过1000条/秒时,消费者收到的JSON事件中部分字段被合并为数组,示例:

{"cs-uri":["tcp://teams.microsoft.com:443/","tcp://teams.microsoft.com:443/"],"cs(Referer)":["-","-"],"dstport":["443","443"],"cs-auth-group":["-","-"],"sc-filter-result":["OBSERVED","OBSERVED"],"@timestamp":"2023-04-18T08:32:52.801Z","srcport":["63917","63917"],"log_type":"KeyValue","sc-bytes":["14767","14767"],"devicetime":["01/12/2022:11:44:19 GMT","01/12/2022:11:44:19 GMT"] ...}

该异常导致消费者无法将事件反序列化为ProxyEvent类,但使用Kafka控制台消费者时事件显示正常。

Java消费者核心配置:

config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
config.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class);
config.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
config.put(JsonDeserializer.VALUE_DEFAULT_TYPE, ProxyEvent.class);
config.put(ConsumerConfig.MAX_PARTITION_FETCH_BYTES_CONFIG, "10485760");
config.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, "10485760");

消费方法:

@KafkaListener(topics = "processed-proxy", groupId = "group_proxy")
public void consume(ProxyEvent proxyEvent) {
     logger.info("event consumed {}", proxyEvent);
     ...
}

事件由Logstash生产者生成,核心输出配置:

output {
    kafka {
        codec => json
        message_key => "%{tenant_uuid}"
        partitioner => "round_robin"
        topic_id => "processed-%{device_type}"
        bootstrap_servers => "${KAFKA_HOST}"
        max_request_size => 10485760
    }
}

问题根源

问题确实与Logstash生产者有关。高吞吐场景下,Logstash的事件处理逻辑可能出现以下情况:

  1. Filter阶段的插件(如mutate/add_field)重复给同一个字段赋值,导致字段被追加为数组;
  2. 批量处理时的并发冲突,导致多个事件的字段被合并到同一个JSON对象中;
  3. 输入插件重复读取数据,导致同一个事件被多次处理,字段被重复设置。

由于控制台消费者是按需消费,大概率未捕获到异常事件,而Java消费者持续高负载消费,更容易触发该问题。

解决方法

1. 排查并修复Logstash事件处理逻辑

  • 替换字段赋值方式:将mutate/add_field改为mutate/replace,避免重复追加字段值:
    mutate {
      replace => {"cs-uri" => "%{source_uri_field}"}
      replace => {"dstport" => "%{destination_port_field}"}
    }
    
  • 调整Pipeline并发配置:根据服务器CPU核心数调整pipeline.workers,同时控制批量处理的大小和延迟,避免高负载下的事件冲突:
    # 在logstash.conf的顶部添加
    pipeline.workers => 4
    pipeline.batch.size => 100
    pipeline.batch.delay => 50
    
  • 检查输入插件配置:确保输入插件(如file、beats)没有重复读取数据的逻辑,比如file插件的start_position设置正确,避免重复消费日志。

2. 调整Logstash Kafka输出插件的批量配置

添加批量发送的控制参数,限制单批发送的事件数量,避免高吞吐下的事件合并:

output {
    kafka {
        codec => json
        message_key => "%{tenant_uuid}"
        partitioner => "round_robin"
        topic_id => "processed-%{device_type}"
        bootstrap_servers => "${KAFKA_HOST}"
        max_request_size => 10485760
        batch_size => 100
        batch_delay => 50
    }
}

3. 为Java消费者添加反序列化容错(临时方案)

如果无法立即修复Logstash配置,可修改ProxyEvent类,添加Jackson注解兼容单值和数组类型:

import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonSetter;
import com.fasterxml.jackson.databind.JsonNode;

public class ProxyEvent {
    // 其他原有字段保留
    
    @JsonProperty("dstport")
    private String dstport;

    @JsonProperty("cs-uri")
    private String csUri;

    @JsonSetter("dstport")
    public void setDstport(Object value) {
        this.setSingleValue(value, "dstport");
    }

    @JsonSetter("cs-uri")
    public void setCsUri(Object value) {
        this.setSingleValue(value, "cs-uri");
    }

    // 通用处理方法,适配其他可能变为数组的字段
    private void setSingleValue(Object value, String fieldName) {
        if (value instanceof JsonNode) {
            JsonNode node = (JsonNode) value;
            if (node.isArray()) {
                this.setFieldValue(fieldName, node.size() > 0 ? node.get(0).asText() : null);
            } else {
                this.setFieldValue(fieldName, node.asText());
            }
        } else {
            this.setFieldValue(fieldName, value.toString());
        }
    }

    private void setFieldValue(String fieldName, String value) {
        switch (fieldName) {
            case "dstport":
                this.dstport = value;
                break;
            case "cs-uri":
                this.csUri = value;
                break;
            // 添加其他字段的case
        }
    }
}

该方式会自动将数组类型的字段取第一个值赋值给字符串字段,避免反序列化失败。

4. 验证Kafka消息内容

使用Kafka控制台消费者持续捕获消息,确认异常事件的来源:

kafka-console-consumer.sh --bootstrap-server ${KAFKA_HOST} --topic processed-proxy --from-beginning > kafka_event_logs.txt

在生成的日志文件中搜索数组格式的字段,进一步定位问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 07:07:10