高吞吐场景下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的事件处理逻辑可能出现以下情况:
- Filter阶段的插件(如
mutate/add_field)重复给同一个字段赋值,导致字段被追加为数组; - 批量处理时的并发冲突,导致多个事件的字段被合并到同一个JSON对象中;
- 输入插件重复读取数据,导致同一个事件被多次处理,字段被重复设置。
由于控制台消费者是按需消费,大概率未捕获到异常事件,而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
相关产品推荐
相关产品推荐

