From ca111610d1259ed6237188deddc132f82c02c47b Mon Sep 17 00:00:00 2001 From: Amritanshu Darbari Date: Tue, 15 Jul 2025 17:50:19 +0530 Subject: [PATCH 1/2] feat: add prometheus/grafana monitoring --- docker-compose.yml | 34 ++ go.mod | 7 + go.sum | 15 + ingestor/main.go | 30 +- monitoring/alertmanager.yml | 13 + .../grafana/dashboards/log-pipeline.json | 478 ++++++++++++++++++ monitoring/grafana/provisioning.yml | 11 + monitoring/prometheus.yml | 20 + monitoring/rules.yml | 12 + parser/main.go | 44 +- storage-writer/main.go | 43 +- 11 files changed, 702 insertions(+), 5 deletions(-) create mode 100644 monitoring/alertmanager.yml create mode 100644 monitoring/grafana/dashboards/log-pipeline.json create mode 100644 monitoring/grafana/provisioning.yml create mode 100644 monitoring/prometheus.yml create mode 100644 monitoring/rules.yml diff --git a/docker-compose.yml b/docker-compose.yml index bb29002..f43d618 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -77,6 +77,8 @@ services: build: context: . dockerfile: parser/Dockerfile + ports: + - "8080" depends_on: kafka: condition: service_healthy @@ -91,6 +93,8 @@ services: build: context: . dockerfile: storage-writer/Dockerfile + ports: + - "8080" depends_on: elasticsearch: condition: service_healthy @@ -100,6 +104,36 @@ services: - KAFKA_BROKERS=kafka:9092 - ELASTICSEARCH_URL=http://elasticsearch:9200 + prometheus: + image: prom/prometheus:v2.53.0 + ports: + - "9090:9090" + volumes: + - ./monitoring/prometheus.yml:/etc/prometheus/prometheus.yml + - ./monitoring/rules.yml:/etc/prometheus/rules.yml + command: + - '--config.file=/etc/prometheus/prometheus.yml' + - '--web.enable-lifecycle' + + alertmanager: + image: prom/alertmanager:v0.27.0 + ports: + - "9093:9093" + volumes: + - ./monitoring/alertmanager.yml:/etc/alertmanager/alertmanager.yml + command: + - '--config.file=/etc/alertmanager/alertmanager.yml' + + grafana: + image: grafana/grafana:11.1.0 + ports: + - "3000:3000" + volumes: + - ./monitoring/grafana/dashboards:/etc/grafana/dashboards + - ./monitoring/grafana/provisioning.yml:/etc/grafana/provisioning/dashboards/provisioning.yml + - grafana_data:/var/lib/grafana + volumes: kafka_data: esdata: + grafana_data: diff --git a/go.mod b/go.mod index 30c0e32..7cbbc1d 100644 --- a/go.mod +++ b/go.mod @@ -10,8 +10,10 @@ require ( ) require ( + github.com/beorn7/perks v1.0.1 // indirect github.com/bytedance/sonic v1.13.3 // indirect github.com/bytedance/sonic/loader v0.2.4 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/cloudwego/base64x v0.1.5 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/eapache/go-resiliency v1.7.0 // indirect @@ -44,8 +46,13 @@ require ( github.com/mattn/go-isatty v0.0.20 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect github.com/modern-go/reflect2 v1.0.2 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/pierrec/lz4/v4 v4.1.22 // indirect + github.com/prometheus/client_golang v1.22.0 // indirect + github.com/prometheus/client_model v0.6.1 // indirect + github.com/prometheus/common v0.62.0 // indirect + github.com/prometheus/procfs v0.15.1 // indirect github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect github.com/sagikazarmark/locafero v0.7.0 // indirect diff --git a/go.sum b/go.sum index c7136f6..93d83d9 100644 --- a/go.sum +++ b/go.sum @@ -1,10 +1,14 @@ github.com/IBM/sarama v1.45.2 h1:8m8LcMCu3REcwpa7fCP6v2fuPuzVwXDAM2DOv3CBrKw= github.com/IBM/sarama v1.45.2/go.mod h1:ppaoTcVdGv186/z6MEKsMm70A5fwJfRTpstI37kVn3Y= +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= github.com/bytedance/sonic v1.13.3 h1:MS8gmaH16Gtirygw7jV91pDCN33NyMrPbN7qiYhEsF0= github.com/bytedance/sonic v1.13.3/go.mod h1:o68xyaF9u2gvVBuGHPlUVCy+ZfmNNO5ETf1+KgkJhz4= github.com/bytedance/sonic/loader v0.1.1/go.mod h1:ncP89zfokxS5LZrJxl5z0UJcsk4M4yY2JpfqGeCtNLU= github.com/bytedance/sonic/loader v0.2.4 h1:ZWCw4stuXUsn1/+zQDqeE7JKP+QO47tz7QCNan80NzY= github.com/bytedance/sonic/loader v0.2.4/go.mod h1:N8A3vUdtUebEY2/VQC0MyhYeKUFosQU6FxH2JmUe6VI= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/cloudwego/base64x v0.1.5 h1:XPciSp1xaq2VCSt6lF0phncD4koWyULpl5bUxbfCyP4= github.com/cloudwego/base64x v0.1.5/go.mod h1:0zlkT4Wn5C6NdauXdJRhSKRlJvmclQ1hhJgA0rcu/8w= github.com/cloudwego/iasm v0.2.0/go.mod h1:8rXZaNYT2n95jn+zTI1sDr+IgcD2GVs0nlbbQPiEFhY= @@ -54,6 +58,7 @@ github.com/golang/snappy v1.0.0 h1:Oy607GVXHs7RtbggtPBnr2RmDArIsAefDwvrdWvRhGs= github.com/golang/snappy v1.0.0/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= @@ -100,12 +105,22 @@ github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4= github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU= github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.22.0 h1:rb93p9lokFEsctTys46VnV1kLCDpVZ0a/Y92Vm0Zc6Q= +github.com/prometheus/client_golang v1.22.0/go.mod h1:R7ljNsLXhuQXYZYtw6GAE9AZg8Y7vEW5scdCXrWRXC0= +github.com/prometheus/client_model v0.6.1 h1:ZKSh/rekM+n3CeS952MLRAdFwIKqeY8b62p8ais2e9E= +github.com/prometheus/client_model v0.6.1/go.mod h1:OrxVMOVHjw3lKMa8+x6HeMGkHMQyHDk9E3jmP2AmGiY= +github.com/prometheus/common v0.62.0 h1:xasJaQlnWAeyHdUBeGjXmutelfJHWMRr+Fg4QszZ2Io= +github.com/prometheus/common v0.62.0/go.mod h1:vyBcEuLSvWos9B1+CyL7JZ2up+uFzXhkqml0W5zIY1I= +github.com/prometheus/procfs v0.15.1 h1:YagwOFzUgYfKKHX6Dr+sHT7km/hxC76UB0learggepc= +github.com/prometheus/procfs v0.15.1/go.mod h1:fB45yRUv8NstnjriLhBQLuOUt+WW4BsoGhij/e3PBqk= github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 h1:bsUq1dX0N8AOIL7EB/X911+m4EHsnWEHeJ0c+3TTBrg= github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= diff --git a/ingestor/main.go b/ingestor/main.go index 66e8449..0d73229 100644 --- a/ingestor/main.go +++ b/ingestor/main.go @@ -14,12 +14,29 @@ import ( "github.com/IBM/sarama" "github.com/MinuteHanD/log-pipeline/config" "github.com/gin-gonic/gin" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" ) const ( MaxMessageSize = 64 * 1024 // max message size 64 kb ) +var ( + logsReceived = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "ingestor_logs_received_total", + Help: "Total number of logs received.", + }, + ) + invalidLogs = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "ingestor_invalid_logs_total", + Help: "Total number of invalid logs.", + }, + ) +) + type LogEntry struct { Timestamp string `json:"timestamp" binding:"required"` Level string `json:"level" binding:"required"` @@ -54,7 +71,7 @@ func NewLogValidator() *LogValidator { "PANIC": true, } - timestampRegex := regexp.MustCompile(`^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d{3})?Z?$`) + timestampRegex := regexp.MustCompile(`^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d{3})?Z?$`) return &LogValidator{ validLevels: validLevels, @@ -177,6 +194,11 @@ func (v *LogValidator) ValidateComplete(data []byte) ValidationResult { } } +func init() { + prometheus.MustRegister(logsReceived) + prometheus.MustRegister(invalidLogs) +} + func main() { logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) @@ -203,6 +225,7 @@ func main() { router := gin.Default() router.POST("/log", func(c *gin.Context) { + logsReceived.Inc() body, err := c.GetRawData() if err != nil { logger.Error("Failed to read request body", "error", err) @@ -215,6 +238,7 @@ func main() { validationResult := validator.ValidateComplete(body) if !validationResult.IsValid { + invalidLogs.Inc() logger.Info("Log validation failed", "validation_errors", validationResult.Errors) c.JSON(http.StatusBadRequest, gin.H{ "error": "validation failed", @@ -261,9 +285,11 @@ func main() { c.JSON(http.StatusOK, gin.H{"status": "ok"}) }) + router.GET("/metrics", gin.WrapH(promhttp.Handler())) + logger.Info("Ingestor service starting", "port", cfg.Ingestor.HTTPPort) if err := router.Run(fmt.Sprintf(":%d", cfg.Ingestor.HTTPPort)); err != nil { logger.Error("Failed to run Gin server", "error", err) os.Exit(1) } -} +} \ No newline at end of file diff --git a/monitoring/alertmanager.yml b/monitoring/alertmanager.yml new file mode 100644 index 0000000..00c31ad --- /dev/null +++ b/monitoring/alertmanager.yml @@ -0,0 +1,13 @@ + +global: + resolve_timeout: 5m + +route: + receiver: 'default-receiver' + group_by: ['alertname', 'cluster', 'service'] + group_wait: 30s + group_interval: 5m + repeat_interval: 4h + +receivers: + - name: 'default-receiver' diff --git a/monitoring/grafana/dashboards/log-pipeline.json b/monitoring/grafana/dashboards/log-pipeline.json new file mode 100644 index 0000000..0df47f1 --- /dev/null +++ b/monitoring/grafana/dashboards/log-pipeline.json @@ -0,0 +1,478 @@ +{ + "annotations": { + "list": [ + { + "builtIn": 1, + "datasource": { + "type": "grafana", + "uid": "-- Grafana --" + }, + "enable": true, + "hide": true, + "iconColor": "rgba(0, 211, 255, 1)", + "name": "Annotations & Alerts", + "type": "dashboard" + } + ] + }, + "editable": true, + "fiscalYearStartMonth": 0, + "graphTooltip": 0, + "id": 1, + "links": [], + "panels": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 0 + }, + "id": 6, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "storage_writer_logs_failed_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Failed to store by storage writer", + "type": "stat" + }, + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 8 + }, + "id": 5, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "storage_writer_logs_stored_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Logs stored by storage writer", + "type": "stat" + }, + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 16 + }, + "id": 4, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "parser_logs_failed_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Logs failed by Parser", + "type": "stat" + }, + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 24 + }, + "id": 3, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "parser_logs_processed_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Logs processed Total (Parser)", + "type": "stat" + }, + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 32 + }, + "id": 2, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "ingestor_invalid_logs_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Invalid Logs (Ingestor)", + "type": "stat" + }, + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "fieldConfig": { + "defaults": { + "color": { + "mode": "thresholds" + }, + "mappings": [], + "thresholds": { + "mode": "absolute", + "steps": [ + { + "color": "green", + "value": null + }, + { + "color": "red", + "value": 80 + } + ] + } + }, + "overrides": [] + }, + "gridPos": { + "h": 8, + "w": 12, + "x": 0, + "y": 40 + }, + "id": 1, + "options": { + "colorMode": "value", + "graphMode": "area", + "justifyMode": "auto", + "orientation": "auto", + "percentChangeColorMode": "standard", + "reduceOptions": { + "calcs": [ + "lastNotNull" + ], + "fields": "", + "values": false + }, + "showPercentChange": false, + "textMode": "auto", + "wideLayout": true + }, + "pluginVersion": "11.1.0", + "targets": [ + { + "datasource": { + "type": "prometheus", + "uid": "ces0211ltt88wc" + }, + "disableTextWrap": false, + "editorMode": "builder", + "expr": "ingestor_logs_received_total", + "fullMetaSearch": false, + "includeNullMetadata": true, + "instant": false, + "legendFormat": "__auto", + "range": true, + "refId": "A", + "useBackend": false + } + ], + "title": "Total Logs recieved", + "type": "stat" + } + ], + "schemaVersion": 39, + "tags": [], + "templating": { + "list": [] + }, + "time": { + "from": "now-6h", + "to": "now" + }, + "timepicker": {}, + "timezone": "browser", + "title": "Log Pipeline Overview", + "uid": "aes035afj4qv4c", + "version": 6, + "weekStart": "" +} \ No newline at end of file diff --git a/monitoring/grafana/provisioning.yml b/monitoring/grafana/provisioning.yml new file mode 100644 index 0000000..725a988 --- /dev/null +++ b/monitoring/grafana/provisioning.yml @@ -0,0 +1,11 @@ +apiVersion: 1 + +providers: + - name: 'default' + orgId: 1 + folder: '' + type: file + disableDeletion: false + editable: true + options: + path: /etc/grafana/dashboards diff --git a/monitoring/prometheus.yml b/monitoring/prometheus.yml new file mode 100644 index 0000000..75a5091 --- /dev/null +++ b/monitoring/prometheus.yml @@ -0,0 +1,20 @@ + +global: + scrape_interval: 15s + +scrape_configs: + - job_name: 'prometheus' + static_configs: + - targets: ['localhost:9090'] + + - job_name: 'ingestor' + static_configs: + - targets: ['ingestor:8081'] + + - job_name: 'parser' + static_configs: + - targets: ['parser:8080'] + + - job_name: 'storage-writer' + static_configs: + - targets: ['storage-writer:8080'] diff --git a/monitoring/rules.yml b/monitoring/rules.yml new file mode 100644 index 0000000..323356a --- /dev/null +++ b/monitoring/rules.yml @@ -0,0 +1,12 @@ + +groups: + - name: example + rules: + - alert: HighErrorRate + expr: rate(errors_total[1m]) > 0.5 + for: 1m + labels: + severity: 'critical' + annotations: + summary: "High error rate detected in {{ $labels.job }}" + description: "The error rate for the {{ $labels.job }} job is over 50% for the last minute." diff --git a/parser/main.go b/parser/main.go index 6c0e962..ba01809 100644 --- a/parser/main.go +++ b/parser/main.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "net/http" "os" "os/signal" "strings" @@ -16,12 +17,35 @@ import ( "github.com/IBM/sarama" "github.com/MinuteHanD/log-pipeline/config" "github.com/MinuteHanD/log-pipeline/kafka" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" ) const ( ConsumerGroup = "parser-group" ) +var ( + logsProcessed = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "parser_logs_processed_total", + Help: "Total number of logs processed.", + }, + ) + logsFailed = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "parser_logs_failed_total", + Help: "Total number of logs that failed to be parsed.", + }, + ) + logsSent = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "parser_logs_sent_total", + Help: "Total number of logs successfully sent to the parsed_logs topic.", + }, + ) +) + type IncomingLogEntry struct { Timestamp string `json:"timestamp"` Level string `json:"level"` @@ -156,7 +180,11 @@ func generateMessageHash(message string) string { return string(rune(hash)) } - +func init() { + prometheus.MustRegister(logsProcessed) + prometheus.MustRegister(logsFailed) + prometheus.MustRegister(logsSent) +} func main() { @@ -168,6 +196,14 @@ func main() { os.Exit(1) } + go func() { + http.Handle("/metrics", promhttp.Handler()) + logger.Info("Starting metrics server on port 8080") + if err := http.ListenAndServe(":8080", nil); err != nil { + logger.Error("Failed to start metrics server", "error", err) + } + }() + config := sarama.NewConfig() config.Version = sarama.V2_8_0_0 config.Consumer.Group.Rebalance.Strategy = sarama.BalanceStrategyRoundRobin @@ -254,10 +290,12 @@ func (h *logHandler) Cleanup(_ sarama.ConsumerGroupSession) error { func (h *logHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for message := range claim.Messages() { + logsProcessed.Inc() h.logger.Info("Received message", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset) var incomingLog IncomingLogEntry if err := json.Unmarshal(message.Value, &incomingLog); err != nil { + logsFailed.Inc() h.logger.Error("Failed to unmarshal incoming log", "error", err, "topic", message.Topic, "offset", message.Offset) kafka.SendToDLQ(h.logger, h.producer, h.cfg.Kafka.Topics.RawDLQ, message, err) @@ -268,6 +306,7 @@ func (h *logHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sar parsedLog, err := h.processor.ProcessLog(incomingLog) if err != nil { + logsFailed.Inc() h.logger.Error("Failed to process log", "error", err, "topic", message.Topic, "offset", message.Offset) kafka.SendToDLQ(h.logger, h.producer, h.cfg.Kafka.Topics.RawDLQ, message, err) @@ -278,6 +317,7 @@ func (h *logHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sar enrichedLogBytes, err := json.Marshal(parsedLog) if err != nil { + logsFailed.Inc() h.logger.Error("Failed to marshal processed log", "error", err, "log_id", parsedLog.ID) session.MarkMessage(message, "") continue @@ -290,8 +330,10 @@ func (h *logHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sar partition, offset, err := h.producer.SendMessage(producerMsg) if err != nil { + logsFailed.Inc() h.logger.Error("Failed to send processed log to Kafka", "error", err, "topic", h.cfg.Kafka.Topics.Parsed, "log_id", parsedLog.ID) } else { + logsSent.Inc() h.logger.Info("Successfully sent processed log to Kafka", "topic", h.cfg.Kafka.Topics.Parsed, "partition", partition, "offset", offset, "log_id", parsedLog.ID) } diff --git a/storage-writer/main.go b/storage-writer/main.go index f736cf8..eaea3b2 100644 --- a/storage-writer/main.go +++ b/storage-writer/main.go @@ -5,6 +5,7 @@ import ( "context" "encoding/json" "fmt" + "net/http" "os" "os/signal" "strings" @@ -19,14 +20,35 @@ import ( "github.com/elastic/go-elasticsearch/v8/esapi" "github.com/MinuteHanD/log-pipeline/config" "github.com/MinuteHanD/log-pipeline/kafka" + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promhttp" ) - - const ( ConsumerGroup = "storage-writer-group" ) +var ( + logsReceived = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "storage_writer_logs_received_total", + Help: "Total number of logs received from the parsed_logs topic.", + }, + ) + logsStored = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "storage_writer_logs_stored_total", + Help: "Total number of logs successfully stored in Elasticsearch.", + }, + ) + logsFailed = prometheus.NewCounter( + prometheus.CounterOpts{ + Name: "storage_writer_logs_failed_total", + Help: "Total number of logs that failed to be stored in Elasticsearch.", + }, + ) +) + type ParsedLog struct { ID string `json:"id"` Timestamp time.Time `json:"timestamp"` @@ -192,6 +214,12 @@ func (lsp *LogStorageProcessor) ProcessAndStore(messageValue []byte) error { return nil } +func init() { + prometheus.MustRegister(logsReceived) + prometheus.MustRegister(logsStored) + prometheus.MustRegister(logsFailed) +} + func main() { logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) @@ -201,6 +229,14 @@ func main() { os.Exit(1) } + go func() { + http.Handle("/metrics", promhttp.Handler()) + logger.Info("Starting metrics server on port 8080") + if err := http.ListenAndServe(":8080", nil); err != nil { + logger.Error("Failed to start metrics server", "error", err) + } + }() + esClient, err := elasticsearch.NewClient(elasticsearch.Config{ Addresses: []string{cfg.Elasticsearch.URL}, }) @@ -317,9 +353,11 @@ func (handler *logStorageHandler) Cleanup(_ sarama.ConsumerGroupSession) error { func (handler *logStorageHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for message := range claim.Messages() { + logsReceived.Inc() handler.logger.Info("Received parsed log", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset) if err := handler.processor.ProcessAndStore(message.Value); err != nil { + logsFailed.Inc() handler.logger.Error("Failed to process and store log", "error", err, "topic", message.Topic, "offset", message.Offset) // Send the failed message to the Dead-Letter Queue @@ -329,6 +367,7 @@ func (handler *logStorageHandler) ConsumeClaim(session sarama.ConsumerGroupSessi continue } + logsStored.Inc() session.MarkMessage(message, "") } From f660dbe39bff0167c6759cb20826efc9e714c8d1 Mon Sep 17 00:00:00 2001 From: Amritanshu Darbari Date: Wed, 16 Jul 2025 15:03:38 +0530 Subject: [PATCH 2/2] fix: Update Readme --- README.md | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 93a0467..1495f44 100644 --- a/README.md +++ b/README.md @@ -159,8 +159,21 @@ The project includes both unit tests and integration tests, which are separated --- +## Monitoring & Observability + +The pipeline exposes key operational metrics from each service, which are collected by Prometheus and visualized in a pre-built Grafana dashboard. This provides at-a-glance visibility into: + +* **Ingestor**: Logs received, invalid logs, and successful handoffs to Kafka. +* **Parser**: Logs processed and logs sent to the Dead-Letter Queue. +* **Storage Writer**: Logs successfully written to Elasticsearch and failures. + +Alerting rules are configured in Alertmanager to notify operators of critical issues, such as a high rate of failed logs or services being down. + +To view the Grafana dashboard, visit `http://localhost:3000` (default credentials: `admin`/`admin`). + +--- + ## Future Work -* **Metrics & Observability**: Expose Prometheus metrics from each service for monitoring and alerting. * **Correlation IDs**: Implement a correlation ID at the `ingestor` and pass it through Kafka headers to trace a single request across all services. * **DLQ Re-processing**: Build a utility or service to consume from the DLQ topics, attempt to re-process messages, and archive unrecoverable ones.