From 1c55781a85a73ccdf93cf27c83367e0c0bfade43 Mon Sep 17 00:00:00 2001 From: nileblack Date: Sat, 7 Nov 2020 21:40:36 +0800 Subject: [PATCH] Catch all exception when consume kafka record (#5760) --- CHANGES.md | 1 + .../agent/kafka/provider/handler/JVMMetricsHandler.java | 5 ++--- .../agent/kafka/provider/handler/MeterServiceHandler.java | 5 ++--- .../agent/kafka/provider/handler/ProfileTaskHandler.java | 5 ++--- .../kafka/provider/handler/ServiceManagementHandler.java | 5 ++--- .../agent/kafka/provider/handler/TraceSegmentHandler.java | 2 +- 6 files changed, 10 insertions(+), 13 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 9d03725e7..e7f0502d1 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -27,6 +27,7 @@ Release Notes. * Improve Kubernetes service registry for ALS analysis. * Add health checker for cluster management * Improve the queryable tags generation. Remove the duplicated tags to reduce the storage payload. +* Fix the threads of the Kafka fetcher exit if some unexpected exceptions happen. * Fix the excessive timeout period set by the kubernetes-client. * Fix deadlock problem when using elasticsearch-client-7.0.0. * Fix storage-jdbc isExists not set dbname. diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/JVMMetricsHandler.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/JVMMetricsHandler.java index c518dfb4e..3fe3ce802 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/JVMMetricsHandler.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/JVMMetricsHandler.java @@ -18,7 +18,6 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler; -import com.google.protobuf.InvalidProtocolBufferException; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.utils.Bytes; @@ -67,8 +66,8 @@ public class JVMMetricsHandler implements KafkaHandler { builder.getMetricsList().forEach(jvmMetric -> { jvmSourceDispatcher.sendMetric(builder.getService(), builder.getServiceInstance(), jvmMetric); }); - } catch (InvalidProtocolBufferException e) { - log.error("", e); + } catch (Exception e) { + log.error("handle record failed", e); } } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/MeterServiceHandler.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/MeterServiceHandler.java index a68c0c500..6a985ecdb 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/MeterServiceHandler.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/MeterServiceHandler.java @@ -18,7 +18,6 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler; -import com.google.protobuf.InvalidProtocolBufferException; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.utils.Bytes; @@ -51,8 +50,8 @@ public class MeterServiceHandler implements KafkaHandler { meterDataCollection.getMeterDataList().forEach(meterData -> processor.read(meterData)); processor.process(); - } catch (InvalidProtocolBufferException e) { - log.error("", e); + } catch (Exception e) { + log.error("handle record failed", e); } } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ProfileTaskHandler.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ProfileTaskHandler.java index 7d0746cca..4d37b3dad 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ProfileTaskHandler.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ProfileTaskHandler.java @@ -18,7 +18,6 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler; -import com.google.protobuf.InvalidProtocolBufferException; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.utils.Bytes; @@ -62,8 +61,8 @@ public class ProfileTaskHandler implements KafkaHandler { snapshotRecord.setTimeBucket(TimeBucket.getRecordTimeBucket(snapshot.getTime())); RecordStreamProcessor.getInstance().in(snapshotRecord); - } catch (InvalidProtocolBufferException e) { - log.error(e.getMessage(), e); + } catch (Exception e) { + log.error("handle record failed", e); } } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandler.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandler.java index ef613f8ea..8f2f5d3bf 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandler.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/ServiceManagementHandler.java @@ -19,7 +19,6 @@ package org.apache.skywalking.oap.server.analyzer.agent.kafka.provider.handler; import com.google.gson.JsonObject; -import com.google.protobuf.InvalidProtocolBufferException; import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; @@ -67,8 +66,8 @@ public class ServiceManagementHandler implements KafkaHandler { } else { keepAlive(InstancePingPkg.parseFrom(record.value().get())); } - } catch (InvalidProtocolBufferException e) { - log.error("", e); + } catch (Exception e) { + log.error("handle record failed", e); } } diff --git a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/TraceSegmentHandler.java b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/TraceSegmentHandler.java index d8bb0ac4b..7a41f5ddf 100644 --- a/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/TraceSegmentHandler.java +++ b/oap-server/server-fetcher-plugin/kafka-fetcher-plugin/src/main/java/org/apache/skywalking/oap/server/analyzer/agent/kafka/provider/handler/TraceSegmentHandler.java @@ -90,7 +90,7 @@ public class TraceSegmentHandler implements KafkaHandler { timer.finish(); } } catch (InvalidProtocolBufferException e) { - log.error(e.getMessage(), e); + log.error("handle record failed", e); } }