From 1184864c14a37aa1c574c38b92340d92b2f0b2db Mon Sep 17 00:00:00 2001 From: Sanket Sudake Date: Tue, 9 Aug 2022 11:06:19 +0530 Subject: [PATCH] Reestablish kakfa consumer group session on disconnection (#2504) * Reestablish kakfa consumer group session on disconnection * Add wait for the consumer * Ignore empty message * Update github.com/Shopify/sarama to v1.35.0 Signed-off-by: Sanket Sudake --- go.mod | 9 +- go.sum | 32 ++- pkg/mqtrigger/messageQueue/kafka/consumer.go | 259 +++++++++++++++++++ pkg/mqtrigger/messageQueue/kafka/kafka.go | 257 ++---------------- pkg/mqtrigger/mqtmanager.go | 2 +- 5 files changed, 306 insertions(+), 253 deletions(-) create mode 100644 pkg/mqtrigger/messageQueue/kafka/consumer.go diff --git a/go.mod b/go.mod index 1ba8208e..2fcaf459 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,7 @@ module github.com/fission/fission go 1.18 require ( - github.com/Shopify/sarama v1.32.0 + github.com/Shopify/sarama v1.35.0 github.com/dchest/uniuri v0.0.0-20200228104902-7aecb25e1fe5 github.com/docopt/docopt-go v0.0.0-20180111231733-ee0de3bc6815 github.com/dustin/go-humanize v1.0.0 @@ -80,7 +80,7 @@ require ( github.com/docker/go-connections v0.4.0 // indirect github.com/docker/go-units v0.4.0 // indirect github.com/dsnet/compress v0.0.2-0.20210315054119-f66993602bf5 // indirect - github.com/eapache/go-resiliency v1.2.0 // indirect + github.com/eapache/go-resiliency v1.3.0 // indirect github.com/eapache/go-xerial-snappy v0.0.0-20180814174437-776d5712da21 // indirect github.com/eapache/queue v1.1.0 // indirect github.com/emirpasic/gods v1.12.0 // indirect @@ -120,7 +120,7 @@ require ( github.com/josharian/intern v1.0.0 // indirect github.com/json-iterator/go v1.1.12 // indirect github.com/kevinburke/ssh_config v0.0.0-20201106050909-4977a11b4351 // indirect - github.com/klauspost/compress v1.14.4 // indirect + github.com/klauspost/compress v1.15.8 // indirect github.com/klauspost/pgzip v1.2.5 // indirect github.com/lib/pq v1.10.4 // indirect github.com/mailru/easyjson v0.7.6 // indirect @@ -136,8 +136,7 @@ require ( github.com/opencontainers/go-digest v1.0.0 // indirect github.com/opencontainers/image-spec v1.0.2 // indirect github.com/opencontainers/runc v1.1.2 // indirect - github.com/pierrec/lz4 v2.6.1+incompatible // indirect - github.com/pierrec/lz4/v4 v4.1.8 // indirect + github.com/pierrec/lz4/v4 v4.1.15 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect github.com/prometheus/client_model v0.2.0 // indirect github.com/prometheus/procfs v0.7.3 // indirect diff --git a/go.sum b/go.sum index 00f28066..34a6cbec 100644 --- a/go.sum +++ b/go.sum @@ -90,10 +90,10 @@ github.com/PuerkitoBio/purell v1.1.1 h1:WEQqlqaGbrPkxLJWfBwQmfEAE1Z7ONdDLqrN38tN github.com/PuerkitoBio/purell v1.1.1/go.mod h1:c11w/QuzBsJSee3cPx9rAFu61PvFxuPbtSwDGJws/X0= github.com/PuerkitoBio/urlesc v0.0.0-20170810143723-de5bf2ad4578 h1:d+Bc7a5rLufV/sSk/8dngufqelfh6jnri85riMAaF/M= github.com/PuerkitoBio/urlesc v0.0.0-20170810143723-de5bf2ad4578/go.mod h1:uGdkoq3SwY9Y+13GIhn11/XLaGBb4BfwItxLd5jeuXE= -github.com/Shopify/sarama v1.32.0 h1:P+RUjEaRU0GMMbYexGMDyrMkLhbbBVUVISDywi+IlFU= -github.com/Shopify/sarama v1.32.0/go.mod h1:+EmJJKZWVT/faR9RcOxJerP+LId4iWdQPBGLy1Y1Njs= -github.com/Shopify/toxiproxy/v2 v2.3.0 h1:62YkpiP4bzdhKMH+6uC5E95y608k3zDwdzuBMsnn3uQ= -github.com/Shopify/toxiproxy/v2 v2.3.0/go.mod h1:KvQTtB6RjCJY4zqNJn7C7JDFgsG5uoHYDirfUfpIm0c= +github.com/Shopify/sarama v1.35.0 h1:opEGHcK8s5OpQF99wW0D4ol7A3qUpfSFigrDXnWmOcs= +github.com/Shopify/sarama v1.35.0/go.mod h1:n8obse6Cz5NjjXjKwR1JeYr7CkQn4KG+HENJ8n/T9oQ= +github.com/Shopify/toxiproxy/v2 v2.4.0 h1:O1e4Jfvr/hefNTNu+8VtdEG5lSeamJRo4aKhMOKNM64= +github.com/Shopify/toxiproxy/v2 v2.4.0/go.mod h1:3ilnjng821bkozDRxNoo64oI/DKqM+rOyJzb564+bvg= github.com/acomagu/bufpipe v1.0.3 h1:fxAGrHZTgQ9w5QqVItgzwj235/uYZYgbXitB+dLupOk= github.com/acomagu/bufpipe v1.0.3/go.mod h1:mxdxdup/WdsKVreO5GpW4+M/1CE2sMG4jeGJ2sYmHc4= github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= @@ -198,8 +198,8 @@ github.com/dsnet/compress v0.0.2-0.20210315054119-f66993602bf5/go.mod h1:qssHWj6 github.com/dsnet/golib v0.0.0-20171103203638-1ea166775780/go.mod h1:Lj+Z9rebOhdfkVLjJ8T6VcRQv3SXugXy999NBtR9aFY= github.com/dustin/go-humanize v1.0.0 h1:VSnTsYCnlFHaM2/igO1h6X3HA71jcobQuxemgkq4zYo= github.com/dustin/go-humanize v1.0.0/go.mod h1:HtrtbFcZ19U5GC7JDqmcUSB87Iq5E25KnS6fMYU6eOk= -github.com/eapache/go-resiliency v1.2.0 h1:v7g92e/KSN71Rq7vSThKaWIq68fL4YHvWyiUKorFR1Q= -github.com/eapache/go-resiliency v1.2.0/go.mod h1:kFI+JgMyC7bLPUVY133qvEBtVayf5mFgVsvEsIPBvNs= +github.com/eapache/go-resiliency v1.3.0 h1:RRL0nge+cWGlxXbUzJ7yMcq6w2XBEr19dCN6HECGaT0= +github.com/eapache/go-resiliency v1.3.0/go.mod h1:5yPzW0MIvSe0JDsv0v+DvcjEv2FyD6iZYSs1ZI+iQho= github.com/eapache/go-xerial-snappy v0.0.0-20180814174437-776d5712da21 h1:YEetp8/yCZMuEPMUDHG0CW/brkkEp8mzqk2+ODEitlw= github.com/eapache/go-xerial-snappy v0.0.0-20180814174437-776d5712da21/go.mod h1:+020luEh2TKB4/GOp8oxxtq0Daoen/Cii55CzbTV6DU= github.com/eapache/queue v1.1.0 h1:YOEu7KNc61ntiQlcEeUIoDTJ2o8mQznoNvUhiigpIqc= @@ -240,8 +240,6 @@ github.com/form3tech-oss/jwt-go v3.2.3+incompatible/go.mod h1:pbq4aXjuKjdthFRnoD github.com/fortytw2/leaktest v1.3.0 h1:u8491cBMTQ8ft8aeV+adlcytMZylmA5nnwwkRZjI8vw= github.com/fortytw2/leaktest v1.3.0/go.mod h1:jDsjWgpAGjm2CA7WthBh/CdZYEPF31XHquHwclZch5g= github.com/frankban/quicktest v1.11.3/go.mod h1:wRf/ReqHper53s+kmmSZizM8NamnL3IM0I9ntUbOk+k= -github.com/frankban/quicktest v1.14.2 h1:SPb1KFFmM+ybpEjPUhCCkZOM5xlovT5UbrMvWnXyBns= -github.com/frankban/quicktest v1.14.2/go.mod h1:mgiwOwqx65TmIk1wJ6Q7wvnVMocbUorkibMOrVTHZps= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ= github.com/fsnotify/fsnotify v1.5.1 h1:mZcQUHVQUQWoPXXtuf9yuEXKudkV2sx1E06UadKWpgI= @@ -369,7 +367,6 @@ github.com/google/go-cmp v0.5.3/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/ github.com/google/go-cmp v0.5.4/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.5.6/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= -github.com/google/go-cmp v0.5.7/go.mod h1:n+brtR0CgQNWTVd5ZUFpTBC8YFBDLK/h/bpaJ8/DtOE= github.com/google/go-cmp v0.5.8 h1:e6P7q2lk1O+qJJb4BtCQXlK8vWEO8V1ZeuEdJNOqZyg= github.com/google/go-cmp v0.5.8/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= @@ -505,8 +502,8 @@ github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.4.1/go.mod h1:RyIbtBH6LamlWaDj8nUwkbUhJ87Yi3uG0guNDohfE1A= github.com/klauspost/compress v1.11.4/go.mod h1:aoV0uJVorq1K+umq18yTdKaF57EivdYsUV+/s2qKfXs= -github.com/klauspost/compress v1.14.4 h1:eijASRJcobkVtSt81Olfh7JX43osYLwy5krOJo6YEu4= -github.com/klauspost/compress v1.14.4/go.mod h1:/3/Vjq9QcHkK5uEr5lBEmyoZ1iFhe47etQ6QUkpK6sk= +github.com/klauspost/compress v1.15.8 h1:JahtItbkWjf2jzm/T+qgMxkP9EMHsqEUA6vCMGmXvhA= +github.com/klauspost/compress v1.15.8/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= github.com/klauspost/cpuid v1.2.0/go.mod h1:Pj4uuM528wm8OyEC2QMXAi2YiTZ96dNQPGgoMS4s3ek= github.com/klauspost/pgzip v1.2.5 h1:qnWYvvKqedOF2ulHpMG72XQol4ILEJ8k2wwRl/Km8oE= github.com/klauspost/pgzip v1.2.5/go.mod h1:Ch1tH69qFZu15pkjo5kYi6mth2Zzwzt50oCQKQE9RUs= @@ -610,11 +607,9 @@ github.com/ory/dockertest v3.3.5+incompatible/go.mod h1:1vX4m9wsvi00u5bseYwXaSnh github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic= github.com/peterbourgon/diskv v2.0.1+incompatible/go.mod h1:uqqh8zWWbv1HBMNONnaR/tNboyR3/BZd58JJSHlUSCU= -github.com/pierrec/lz4 v2.6.1+incompatible h1:9UY3+iC23yxF0UfGaYrGplQ+79Rg+h/q9FV9ix19jjM= -github.com/pierrec/lz4 v2.6.1+incompatible/go.mod h1:pdkljMzZIN41W+lC3N2tnIh5sFi+IEE17M5jbnwPHcY= github.com/pierrec/lz4/v4 v4.1.2/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= -github.com/pierrec/lz4/v4 v4.1.8 h1:ieHkV+i2BRzngO4Wd/3HGowuZStgq6QkPsD1eolNAO4= -github.com/pierrec/lz4/v4 v4.1.8/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -733,8 +728,8 @@ github.com/wcharczuk/go-chart v2.0.1+incompatible/go.mod h1:PF5tmL4EIx/7Wf+hEkpC github.com/xanzy/ssh-agent v0.3.0 h1:wUMzuKtKilRgBAD1sUb8gOwwRr2FGoBVumcjoOACClI= github.com/xanzy/ssh-agent v0.3.0/go.mod h1:3s9xbODqPuuhK9JV1R321M/FlMZSBvE5aY6eAcqrDh0= github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= -github.com/xdg-go/scram v1.1.0/go.mod h1:1WAq6h33pAW+iRreB34OORO2Nf7qel3VV3fjBj+hCSs= -github.com/xdg-go/stringprep v1.0.2/go.mod h1:8F9zXuvzgwmyT5DUm4GUfZGDdT3W+LCvS6+da4O5kxM= +github.com/xdg-go/scram v1.1.1/go.mod h1:RaEWvsqvNKKvBPvcKeFjrG2cJqOkHTiyTpzz23ni57g= +github.com/xdg-go/stringprep v1.0.3/go.mod h1:W3f5j4i+9rC0kuIEJL0ky1VpHXQU3ocBgklLGvcBnW8= github.com/xi2/xz v0.0.0-20171230120015-48954b6210f8 h1:nIPpBwaJSVYIxUFsDv3M8ofmx9yWTog9BfvIu0q41lo= github.com/xi2/xz v0.0.0-20171230120015-48954b6210f8/go.mod h1:HUYIGzjTL3rfEspMxjDjgmT5uz5wzYJKVo23qUhYTos= github.com/xiang90/probing v0.0.0-20190116061207-43a291ad63a2/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU= @@ -931,6 +926,7 @@ golang.org/x/net v0.0.0-20211015210444-4f30a5c0130f/go.mod h1:9nx3DQGgdP8bBQD5qx golang.org/x/net v0.0.0-20211112202133-69e39bad7dc2/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= golang.org/x/net v0.0.0-20211216030914-fe4d6282115f/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= golang.org/x/net v0.0.0-20220127200216-cd36cc0744dd/go.mod h1:CfG3xpIq0wQ8r1q4Su4UZFWDARRcnwPjda9FqA0JpMk= +golang.org/x/net v0.0.0-20220708220712-1185a9018129/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= golang.org/x/net v0.0.0-20220805013720-a33c5aa5df48 h1:N9Vc/rorQUDes6B9CNdIxAn5jODGj2wzfrei2x4wNj4= golang.org/x/net v0.0.0-20220805013720-a33c5aa5df48/go.mod h1:YDH+HFinaLZZlnHAfSS6ZXJJ9M9t4Dl22yv3iI2vPwk= golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= @@ -961,6 +957,7 @@ golang.org/x/sync v0.0.0-20200317015054-43a5402ce75a/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20210220032951-036812b2e83c h1:5KslGYwFpkhGh+Q16bwMP3cOontH8FOep7tGV86Y7SQ= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sys v0.0.0-20180823144017-11551d06cbcc/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -1045,6 +1042,7 @@ golang.org/x/sys v0.0.0-20211124211545-fe61309f8881/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220114195835-da31bd327af9/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220209214540-3681064d5158/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220728004956-3c1f35247d10 h1:WIoqL4EROvwiPdUtaip4VcDdpZ4kha7wBWZrbVKCIZg= golang.org/x/sys v0.0.0-20220728004956-3c1f35247d10/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= diff --git a/pkg/mqtrigger/messageQueue/kafka/consumer.go b/pkg/mqtrigger/messageQueue/kafka/consumer.go new file mode 100644 index 00000000..26b31968 --- /dev/null +++ b/pkg/mqtrigger/messageQueue/kafka/consumer.go @@ -0,0 +1,259 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package kafka + +import ( + "fmt" + "io" + "net/http" + "strconv" + "strings" + + "github.com/Shopify/sarama" + "github.com/pkg/errors" + "go.uber.org/zap" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger" + "github.com/fission/fission/pkg/utils" +) + +type MqtConsumerGroupHandler struct { + version sarama.KafkaVersion + logger *zap.Logger + trigger *fv1.MessageQueueTrigger + fissionHeaders map[string]string + producer sarama.SyncProducer + fnUrl string + ready chan bool +} + +func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, + logger *zap.Logger, + trigger *fv1.MessageQueueTrigger, + producer sarama.SyncProducer, + routerUrl string) MqtConsumerGroupHandler { + ch := MqtConsumerGroupHandler{ + version: version, + logger: logger, + trigger: trigger, + producer: producer, + ready: make(chan bool), + } + // Support other function ref types + if ch.trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { + ch.logger.Fatal("unsupported function reference type for trigger", + zap.Any("function_reference_type", ch.trigger.Spec.FunctionReference.Type), + zap.String("trigger", ch.trigger.ObjectMeta.Name)) + } + // Generate the Headers + ch.fissionHeaders = map[string]string{ + "X-Fission-MQTrigger-Topic": ch.trigger.Spec.Topic, + "X-Fission-MQTrigger-RespTopic": ch.trigger.Spec.ResponseTopic, + "X-Fission-MQTrigger-ErrorTopic": ch.trigger.Spec.ErrorTopic, + "Content-Type": ch.trigger.Spec.ContentType, + } + ch.fnUrl = routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(ch.trigger.Spec.FunctionReference.Name, ch.trigger.ObjectMeta.Namespace), "/") + ch.logger.Debug("function HTTP URL", zap.String("url", ch.fnUrl)) + return ch +} + +// Setup implemented to satisfy the sarama.ConsumerGroupHandler interface +func (ch MqtConsumerGroupHandler) Setup(session sarama.ConsumerGroupSession) error { + ch.logger.With( + zap.String("trigger", ch.trigger.ObjectMeta.Name), + zap.String("topic", ch.trigger.Spec.Topic), + zap.String("memberID", session.MemberID()), + zap.Int32("generationID", session.GenerationID()), + zap.String("claims", fmt.Sprintf("%v", session.Claims())), + ).Info("consumer group session setup") + // Mark the consumer as ready + close(ch.ready) + return nil +} + +// Cleanup implemented to satisfy the sarama.ConsumerGroupHandler interface +func (ch MqtConsumerGroupHandler) Cleanup(session sarama.ConsumerGroupSession) error { + ch.logger.With( + zap.String("trigger", ch.trigger.ObjectMeta.Name), + zap.String("topic", ch.trigger.Spec.Topic), + zap.String("memberID", session.MemberID()), + zap.Int32("generationID", session.GenerationID()), + zap.String("claims", fmt.Sprintf("%v", session.Claims())), + ).Info("consumer group session cleanup") + return nil +} + +// ConsumeClaims implemented to satisfy the sarama.ConsumerGroupHandler interface +func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { + // Do not move the code below to a goroutine. + // The `ConsumeClaim` itself is called within a goroutine + for { + select { + case msg := <-claim.Messages(): + if msg != nil { + ch.kafkaMsgHandler(msg) + session.MarkMessage(msg, "") + mqtrigger.IncreaseMessageCount(ch.trigger.Name, ch.trigger.Namespace) + } + // Should return when `session.Context()` is done. + case <-session.Context().Done(): + return nil + } + } +} + +func (ch *MqtConsumerGroupHandler) kafkaMsgHandler(msg *sarama.ConsumerMessage) { + value := string(msg.Value) + + // Create request + req, err := http.NewRequest("POST", ch.fnUrl, strings.NewReader(value)) + if err != nil { + ch.logger.Error("failed to create HTTP request to invoke function", + zap.Error(err), + zap.String("function_url", ch.fnUrl)) + return + } + + // Set the headers came from Kafka record + // Using Header.Add() as msg.Headers may have keys with more than one value + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + for _, h := range msg.Headers { + req.Header.Add(string(h.Key), string(h.Value)) + } + } else { + ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", + zap.Any("current_version", ch.version)) + } + + for k, v := range ch.fissionHeaders { + req.Header.Set(k, v) + } + + // Make the request + var resp *http.Response + for attempt := 0; attempt <= ch.trigger.Spec.MaxRetries; attempt++ { + // Make the request + resp, err = http.DefaultClient.Do(req) + if err != nil { + ch.logger.Error("sending function invocation request failed", + zap.Error(err), + zap.String("function_url", ch.fnUrl), + zap.String("trigger", ch.trigger.ObjectMeta.Name)) + continue + } + if resp == nil { + continue + } + if err == nil && resp.StatusCode == http.StatusOK { + // Success, quit retrying + break + } + } + + generateErrorHeaders := func(errString string) []sarama.RecordHeader { + var errorHeaders []sarama.RecordHeader + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + if count, ok := errorMessageMap[errString]; ok { + errorMessageMap[errString] = count + 1 + } else { + errorMessageMap[errString] = 1 + } + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("MessageSource"), Value: []byte(ch.trigger.Spec.Topic)}) + errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("RecycleCounter"), Value: []byte(strconv.Itoa(errorMessageMap[errString]))}) + } + return errorHeaders + } + + if resp == nil { + errorString := fmt.Sprintf("request exceed retries: %v", ch.trigger.Spec.MaxRetries) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + fmt.Errorf(errorString), errorHeaders) + return + } + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + + ch.logger.Debug("got response from function invocation", + zap.String("function_url", ch.fnUrl), + zap.String("trigger", ch.trigger.ObjectMeta.Name), + zap.String("body", string(body))) + + if err != nil { + errorString := "request body error: " + string(body) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + errors.Wrapf(err, errorString), errorHeaders) + return + } + if resp.StatusCode != 200 { + errorString := fmt.Sprintf("request returned failure: %v, request body error: %v", resp.StatusCode, body) + errorHeaders := generateErrorHeaders(errorString) + errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, + fmt.Errorf("request returned failure: %v", resp.StatusCode), errorHeaders) + return + } + if len(ch.trigger.Spec.ResponseTopic) > 0 { + // Generate Kafka record headers + var kafkaRecordHeaders []sarama.RecordHeader + if ch.version.IsAtLeast(sarama.V0_11_0_0) { + for k, v := range resp.Header { + // One key may have multiple values + for _, v := range v { + kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) + } + } + } else { + ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", + zap.Any("current_version", ch.version)) + } + + _, _, err := ch.producer.SendMessage(&sarama.ProducerMessage{ + Topic: ch.trigger.Spec.ResponseTopic, + Value: sarama.StringEncoder(body), + Headers: kafkaRecordHeaders, + }) + if err != nil { + ch.logger.Warn("failed to publish response body from function invocation to topic", + zap.Error(err), + zap.String("topic", ch.trigger.Spec.Topic), + zap.String("function_url", ch.fnUrl)) + return + } + } +} + +func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) { + if len(trigger.Spec.ErrorTopic) > 0 { + _, _, e := producer.SendMessage(&sarama.ProducerMessage{ + Topic: trigger.Spec.ErrorTopic, + Value: sarama.StringEncoder(err.Error()), + Headers: errorTopicHeaders, + }) + if e != nil { + logger.Error("failed to publish message to error topic", + zap.Error(e), + zap.String("trigger", trigger.ObjectMeta.Name), + zap.String("message", err.Error()), + zap.String("topic", trigger.Spec.Topic)) + } + } else { + logger.Error("message received to publish to error topic, but no error topic was set", + zap.String("message", err.Error()), zap.String("trigger", trigger.ObjectMeta.Name), zap.String("function_url", funcUrl)) + } +} diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index 8cc8bf1b..69e66551 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -20,9 +20,6 @@ import ( "context" "crypto/tls" "crypto/x509" - "fmt" - "io" - "net/http" "os" "regexp" "strconv" @@ -33,11 +30,9 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" - "github.com/fission/fission/pkg/mqtrigger" "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/mqtrigger/validator" - "github.com/fission/fission/pkg/utils" ) func init() { @@ -66,190 +61,12 @@ type ( Factory struct{} ) -type MqtConsumerGroupHandler struct { - version sarama.KafkaVersion - logger *zap.Logger - trigger *fv1.MessageQueueTrigger - fissionHeaders map[string]string - producer sarama.SyncProducer - fnUrl string -} - type MqtConsumer struct { ctx context.Context cancel context.CancelFunc consumer sarama.ConsumerGroup } -func NewMqtConsumerGroupHandler(version sarama.KafkaVersion, - logger *zap.Logger, - trigger *fv1.MessageQueueTrigger, - producer sarama.SyncProducer, - routerUrl string) MqtConsumerGroupHandler { - ch := MqtConsumerGroupHandler{ - version: version, - logger: logger, - trigger: trigger, - producer: producer, - } - // Support other function ref types - if ch.trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { - ch.logger.Fatal("unsupported function reference type for trigger", - zap.Any("function_reference_type", ch.trigger.Spec.FunctionReference.Type), - zap.String("trigger", ch.trigger.ObjectMeta.Name)) - } - // Generate the Headers - ch.fissionHeaders = map[string]string{ - "X-Fission-MQTrigger-Topic": ch.trigger.Spec.Topic, - "X-Fission-MQTrigger-RespTopic": ch.trigger.Spec.ResponseTopic, - "X-Fission-MQTrigger-ErrorTopic": ch.trigger.Spec.ErrorTopic, - "Content-Type": ch.trigger.Spec.ContentType, - } - ch.fnUrl = routerUrl + "/" + strings.TrimPrefix(utils.UrlForFunction(ch.trigger.Spec.FunctionReference.Name, ch.trigger.ObjectMeta.Namespace), "/") - ch.logger.Debug("function HTTP URL", zap.String("url", ch.fnUrl)) - return ch -} - -func (ch MqtConsumerGroupHandler) Setup(sarama.ConsumerGroupSession) error { - return nil -} - -func (ch MqtConsumerGroupHandler) Cleanup(sarama.ConsumerGroupSession) error { - return nil -} - -func (ch MqtConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { - for msg := range claim.Messages() { - ch.kafkaMsgHandler(session, msg) - mqtrigger.IncreaseMessageCount(ch.trigger.Name, ch.trigger.Namespace) - } - return nil -} - -//func (ch *MqtConsumerGroupHandler) kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *fv1.MessageQueueTrigger, msg *sarama.ConsumerMessage, consumer *cluster.Consumer) { - -func (ch *MqtConsumerGroupHandler) kafkaMsgHandler(session sarama.ConsumerGroupSession, msg *sarama.ConsumerMessage) { - var value string = string(msg.Value[:]) - - // Create request - req, err := http.NewRequest("POST", ch.fnUrl, strings.NewReader(value)) - if err != nil { - ch.logger.Error("failed to create HTTP request to invoke function", - zap.Error(err), - zap.String("function_url", ch.fnUrl)) - return - } - - // Set the headers came from Kafka record - // Using Header.Add() as msg.Headers may have keys with more than one value - if ch.version.IsAtLeast(sarama.V0_11_0_0) { - for _, h := range msg.Headers { - req.Header.Add(string(h.Key), string(h.Value)) - } - } else { - ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", - zap.Any("current_version", ch.version)) - } - - for k, v := range ch.fissionHeaders { - req.Header.Set(k, v) - } - - // Make the request - var resp *http.Response - for attempt := 0; attempt <= ch.trigger.Spec.MaxRetries; attempt++ { - // Make the request - resp, err = http.DefaultClient.Do(req) - if err != nil { - ch.logger.Error("sending function invocation request failed", - zap.Error(err), - zap.String("function_url", ch.fnUrl), - zap.String("trigger", ch.trigger.ObjectMeta.Name)) - continue - } - if resp == nil { - continue - } - if err == nil && resp.StatusCode == http.StatusOK { - // Success, quit retrying - break - } - } - - generateErrorHeaders := func(errString string) []sarama.RecordHeader { - var errorHeaders []sarama.RecordHeader - if ch.version.IsAtLeast(sarama.V0_11_0_0) { - if count, ok := errorMessageMap[errString]; ok { - errorMessageMap[errString] = count + 1 - } else { - errorMessageMap[errString] = 1 - } - errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("MessageSource"), Value: []byte(ch.trigger.Spec.Topic)}) - errorHeaders = append(errorHeaders, sarama.RecordHeader{Key: []byte("RecycleCounter"), Value: []byte(strconv.Itoa(errorMessageMap[errString]))}) - } - return errorHeaders - } - - if resp == nil { - errorString := fmt.Sprintf("request exceed retries: %v", ch.trigger.Spec.MaxRetries) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, - fmt.Errorf(errorString), errorHeaders) - return - } - defer resp.Body.Close() - body, err := io.ReadAll(resp.Body) - - ch.logger.Debug("got response from function invocation", - zap.String("function_url", ch.fnUrl), - zap.String("trigger", ch.trigger.ObjectMeta.Name), - zap.String("body", string(body))) - - if err != nil { - errorString := "request body error: " + string(body) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, - errors.Wrapf(err, errorString), errorHeaders) - return - } - if resp.StatusCode != 200 { - errorString := fmt.Sprintf("request returned failure: %v, request body error: %v", resp.StatusCode, body) - errorHeaders := generateErrorHeaders(errorString) - errorHandler(ch.logger, ch.trigger, ch.producer, ch.fnUrl, - fmt.Errorf("request returned failure: %v", resp.StatusCode), errorHeaders) - return - } - if len(ch.trigger.Spec.ResponseTopic) > 0 { - // Generate Kafka record headers - var kafkaRecordHeaders []sarama.RecordHeader - if ch.version.IsAtLeast(sarama.V0_11_0_0) { - for k, v := range resp.Header { - // One key may have multiple values - for _, v := range v { - kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) - } - } - } else { - ch.logger.Warn("headers are not supported by current Kafka version, needs v0.11+: no record headers to add in HTTP request", - zap.Any("current_version", ch.version)) - } - - _, _, err := ch.producer.SendMessage(&sarama.ProducerMessage{ - Topic: ch.trigger.Spec.ResponseTopic, - Value: sarama.StringEncoder(body), - Headers: kafkaRecordHeaders, - }) - if err != nil { - ch.logger.Warn("failed to publish response body from function invocation to topic", - zap.Error(err), - zap.String("topic", ch.trigger.Spec.Topic), - zap.String("function_url", ch.fnUrl)) - return - } - } - session.MarkMessage(msg, "") -} - func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { return New(logger, mqCfg, routerUrl) } @@ -297,8 +114,8 @@ func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messa } func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { - kafka.logger.Info("inside kakfa subscribe", zap.Any("trigger", trigger)) - kafka.logger.Info("brokers set", zap.Strings("brokers", kafka.brokers)) + kafka.logger.Debug("inside kakfa subscribe", zap.Any("trigger", trigger)) + kafka.logger.Debug("brokers set", zap.Strings("brokers", kafka.brokers)) // Create new consumer consumerConfig := sarama.NewConfig() @@ -314,48 +131,40 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub // Setup TLS for both producer and consumer if kafka.tls { - consumerConfig.Net.TLS.Enable = true - producerConfig.Net.TLS.Enable = true tlsConfig, err := kafka.getTLSConfig() if err != nil { return nil, err } + producerConfig.Net.TLS.Enable = true producerConfig.Net.TLS.Config = tlsConfig + consumerConfig.Net.TLS.Enable = true consumerConfig.Net.TLS.Config = tlsConfig } consumer, err := sarama.NewConsumerGroup(kafka.brokers, string(trigger.ObjectMeta.UID), consumerConfig) - // consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.ObjectMeta.UID), []string{trigger.Spec.Topic}, consumerConfig) - kafka.logger.Info("created a new consumer", zap.Strings("brokers", kafka.brokers), - zap.String("input topic", trigger.Spec.Topic), - zap.String("output topic", trigger.Spec.ResponseTopic), - zap.String("error topic", trigger.Spec.ErrorTopic), - zap.String("trigger name", trigger.ObjectMeta.Name), - zap.String("function namespace", trigger.ObjectMeta.Namespace), - zap.String("function name", trigger.Spec.FunctionReference.Name)) if err != nil { return nil, err } producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig) - kafka.logger.Info("created a new producer", zap.Strings("brokers", kafka.brokers), - zap.String("input topic", trigger.Spec.Topic), - zap.String("output topic", trigger.Spec.ResponseTopic), - zap.String("error topic", trigger.Spec.ErrorTopic), - zap.String("trigger name", trigger.ObjectMeta.Name), - zap.String("function namespace", trigger.ObjectMeta.Namespace), - zap.String("function name", trigger.Spec.FunctionReference.Name)) - if err != nil { return nil, err } + kafka.logger.Info("created a new producer and a new consumer", zap.Strings("brokers", kafka.brokers), + zap.String("topic", trigger.Spec.Topic), + zap.String("response topic", trigger.Spec.ResponseTopic), + zap.String("error topic", trigger.Spec.ErrorTopic), + zap.String("trigger", trigger.ObjectMeta.Name), + zap.String("function namespace", trigger.ObjectMeta.Namespace), + zap.String("function name", trigger.Spec.FunctionReference.Name)) + // consume errors go func() { for err := range consumer.Errors() { - kafka.logger.Error("consumer error", zap.Error(err)) + kafka.logger.With(zap.String("trigger", trigger.ObjectMeta.Name), zap.String("topic", trigger.Spec.Topic)).Error("consumer error received", zap.Error(err)) } }() @@ -365,16 +174,24 @@ func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Sub // consume messages go func() { topic := []string{trigger.Spec.Topic} - err = consumer.Consume(ctx, topic, ch) - if err != nil { - kafka.logger.Error("consumer error", zap.Error(err)) - } + // Create a new session for the consumer group until the context is cancelled + for { + // Consume messages + err := consumer.Consume(ctx, topic, ch) + if err != nil { + kafka.logger.Error("consumer error", zap.Error(err), zap.String("trigger", trigger.ObjectMeta.Name)) + } - if ctx.Err() != nil { - return + if ctx.Err() != nil { + kafka.logger.Info("consumer context cancelled", zap.String("trigger", trigger.ObjectMeta.Name)) + return + } + ch.ready = make(chan bool) } }() + <-ch.ready // wait for consumer to be ready + mqtConsumer := MqtConsumer{ ctx: ctx, cancel: cancel, @@ -413,26 +230,6 @@ func (kafka Kafka) Unsubscribe(subscription messageQueue.Subscription) error { return mqtConsumer.consumer.Close() } -func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer sarama.SyncProducer, funcUrl string, err error, errorTopicHeaders []sarama.RecordHeader) { - if len(trigger.Spec.ErrorTopic) > 0 { - _, _, e := producer.SendMessage(&sarama.ProducerMessage{ - Topic: trigger.Spec.ErrorTopic, - Value: sarama.StringEncoder(err.Error()), - Headers: errorTopicHeaders, - }) - if e != nil { - logger.Error("failed to publish message to error topic", - zap.Error(e), - zap.String("trigger", trigger.ObjectMeta.Name), - zap.String("message", err.Error()), - zap.String("topic", trigger.Spec.Topic)) - } - } else { - logger.Error("message received to publish to error topic, but no error topic was set", - zap.String("message", err.Error()), zap.String("trigger", trigger.ObjectMeta.Name), zap.String("function_url", funcUrl)) - } -} - // The validation is based on Kafka's internal implementation: // https://github.com/apache/kafka/blob/cde6d18983b5d58199f8857d8d61d7efcbe6e54a/clients/src/main/java/org/apache/kafka/common/internals/Topic.java#L36-L47 func IsTopicValid(topic string) bool { diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 77c9df18..b416f63b 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -155,7 +155,7 @@ func (mqt *MessageQueueTriggerManager) delTriggerSubscription(trigger *fv1.Messa func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueueTrigger) { isPresent := mqt.checkTriggerSubscription(trigger) if isPresent { - mqt.logger.Info("message queue trigger already registered", zap.String("trigger_name", trigger.ObjectMeta.Name)) + mqt.logger.Debug("message queue trigger already registered", zap.String("trigger_name", trigger.ObjectMeta.Name)) return }