Skip to content

Commit

Permalink
Add dcp flow control (#18)
Browse files Browse the repository at this point in the history
* Add dcp flow control

* Fix single select case warning

* Fix go lint errors
  • Loading branch information
mhmtszr authored Mar 9, 2023
1 parent 977221f commit ad5abb6
Show file tree
Hide file tree
Showing 5 changed files with 291 additions and 277 deletions.
20 changes: 9 additions & 11 deletions connector.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package gokafkaconnectcouchbase

import (
"github.com/Trendyol/go-dcp-client/models"
"github.com/Trendyol/go-kafka-connect-couchbase/kafka/message"

godcpclient "github.com/Trendyol/go-dcp-client"
Expand Down Expand Up @@ -36,19 +37,16 @@ func (c *connector) Close() {
}
}

func (c *connector) listener(event interface{}, err error) {
if err != nil {
c.errorLogger.Printf("error | %v", err)
return
}
const defaultCollection = "_default"

func (c *connector) listener(ctx *models.ListenerContext) {
var e couchbase.Event
switch event := event.(type) {
case godcpclient.DcpMutation:
switch event := ctx.Event.(type) {
case models.DcpMutation:
e = couchbase.NewMutateEvent(event.Key, event.Value, event.CollectionName)
case godcpclient.DcpExpiration:
case models.DcpExpiration:
e = couchbase.NewExpireEvent(event.Key, nil, event.CollectionName)
case godcpclient.DcpDeletion:
case models.DcpDeletion:
e = couchbase.NewDeleteEvent(event.Key, nil, event.CollectionName)
default:
return
Expand All @@ -58,7 +56,7 @@ func (c *connector) listener(event interface{}, err error) {
defer message.MessagePool.Put(kafkaMessage)
var collectionName string
if e.CollectionName == nil {
collectionName = "_default"
collectionName = defaultCollection
} else {
collectionName = *e.CollectionName
}
Expand All @@ -67,7 +65,7 @@ func (c *connector) listener(event interface{}, err error) {
c.errorLogger.Printf("unexpected collection | %s", collectionName)
return
}
c.producer.Produce(kafkaMessage.Value, kafkaMessage.Key, kafkaMessage.Headers, topic)
c.producer.Produce(ctx, kafkaMessage.Value, kafkaMessage.Key, kafkaMessage.Headers, topic)
}
}

Expand Down
96 changes: 58 additions & 38 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,71 +1,91 @@
module github.com/Trendyol/go-kafka-connect-couchbase

go 1.19
go 1.20

require (
github.com/Trendyol/go-dcp-client v0.0.8
github.com/gookit/config/v2 v2.1.8
github.com/segmentio/kafka-go v0.4.38
github.com/Trendyol/go-dcp-client v0.0.14
github.com/gookit/config/v2 v2.2.1
github.com/segmentio/kafka-go v0.4.39
)

require (
github.com/andybalholm/brotli v1.0.4 // indirect
github.com/ansrivas/fiberprometheus/v2 v2.4.1 // indirect
github.com/avast/retry-go/v4 v4.3.1 // indirect
github.com/VividCortex/ewma v1.2.0 // indirect
github.com/andybalholm/brotli v1.0.5 // indirect
github.com/ansrivas/fiberprometheus/v2 v2.6.0 // indirect
github.com/avast/retry-go/v4 v4.3.3 // indirect
github.com/beorn7/perks v1.0.1 // indirect
github.com/cespare/xxhash/v2 v2.1.2 // indirect
github.com/couchbase/gocbcore/v10 v10.2.0 // indirect
github.com/cespare/xxhash/v2 v2.2.0 // indirect
github.com/couchbase/gocbcore/v10 v10.2.1 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/go-logr/logr v1.2.2 // indirect
github.com/gofiber/adaptor/v2 v2.1.25 // indirect
github.com/gofiber/fiber/v2 v2.39.0 // indirect
github.com/emicklei/go-restful/v3 v3.10.1 // indirect
github.com/fatih/color v1.14.1 // indirect
github.com/go-logr/logr v1.2.3 // indirect
github.com/go-openapi/jsonpointer v0.19.6 // indirect
github.com/go-openapi/jsonreference v0.20.2 // indirect
github.com/go-openapi/swag v0.22.3 // indirect
github.com/goccy/go-yaml v1.10.0 // indirect
github.com/gofiber/adaptor/v2 v2.1.32 // indirect
github.com/gofiber/fiber/v2 v2.42.0 // indirect
github.com/gogo/protobuf v1.3.2 // indirect
github.com/golang/protobuf v1.5.2 // indirect
github.com/golang/snappy v0.0.4 // indirect
github.com/google/go-cmp v0.5.8 // indirect
github.com/google/gnostic v0.6.9 // indirect
github.com/google/go-cmp v0.5.9 // indirect
github.com/google/gofuzz v1.2.0 // indirect
github.com/google/uuid v1.3.0 // indirect
github.com/googleapis/gnostic v0.5.5 // indirect
github.com/gookit/goutil v0.5.15 // indirect
github.com/gookit/color v1.5.2 // indirect
github.com/gookit/goutil v0.6.6 // indirect
github.com/imdario/mergo v0.3.13 // indirect
github.com/josharian/intern v1.0.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/compress v1.15.9 // indirect
github.com/klauspost/compress v1.16.0 // indirect
github.com/mailru/easyjson v0.7.7 // indirect
github.com/mattn/go-colorable v0.1.13 // indirect
github.com/mattn/go-isatty v0.0.16 // indirect
github.com/mattn/go-isatty v0.0.17 // indirect
github.com/mattn/go-runewidth v0.0.14 // indirect
github.com/matttproud/golang_protobuf_extensions v1.0.2-0.20181231171920-c182affec369 // indirect
github.com/matttproud/golang_protobuf_extensions v1.0.4 // indirect
github.com/mitchellh/mapstructure v1.5.0 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/pierrec/lz4/v4 v4.1.15 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/philhofer/fwd v1.1.2 // indirect
github.com/pierrec/lz4/v4 v4.1.17 // indirect
github.com/prometheus/client_golang v1.14.0 // indirect
github.com/prometheus/client_model v0.3.0 // indirect
github.com/prometheus/common v0.37.0 // indirect
github.com/prometheus/procfs v0.8.0 // indirect
github.com/rivo/uniseg v0.2.0 // indirect
github.com/rs/zerolog v1.28.0 // indirect
github.com/prometheus/common v0.42.0 // indirect
github.com/prometheus/procfs v0.9.0 // indirect
github.com/rivo/uniseg v0.4.4 // indirect
github.com/rs/zerolog v1.29.0 // indirect
github.com/savsgio/dictpool v0.0.0-20221023140959-7bf2e61cea94 // indirect
github.com/savsgio/gotils v0.0.0-20230208104028-c358bd845dee // indirect
github.com/tinylib/msgp v1.1.8 // indirect
github.com/valyala/bytebufferpool v1.0.0 // indirect
github.com/valyala/fasthttp v1.40.0 // indirect
github.com/valyala/fasthttp v1.44.0 // indirect
github.com/valyala/tcplisten v1.0.0 // indirect
github.com/xdg/scram v1.0.5 // indirect
github.com/xdg/stringprep v1.0.3 // indirect
golang.org/x/crypto v0.0.0-20220829220503-c86fa9a7ed90 // indirect
golang.org/x/net v0.0.0-20220722155237-a158d28d115b // indirect
golang.org/x/oauth2 v0.0.0-20220223155221-ee480838109b // indirect
golang.org/x/sys v0.0.0-20220829200755-d48e67d00261 // indirect
golang.org/x/term v0.0.0-20220722155259-a9ba230a4035 // indirect
golang.org/x/text v0.3.8 // indirect
golang.org/x/time v0.0.0-20210723032227-1f47c861a9ac // indirect
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect
golang.org/x/crypto v0.7.0 // indirect
golang.org/x/net v0.8.0 // indirect
golang.org/x/oauth2 v0.6.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/sys v0.6.0 // indirect
golang.org/x/term v0.6.0 // indirect
golang.org/x/text v0.8.0 // indirect
golang.org/x/time v0.3.0 // indirect
golang.org/x/xerrors v0.0.0-20220907171357-04be3eba64a2 // indirect
google.golang.org/appengine v1.6.7 // indirect
google.golang.org/protobuf v1.28.1 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
k8s.io/api v0.22.5 // indirect
k8s.io/apimachinery v0.22.5 // indirect
k8s.io/client-go v0.22.5 // indirect
k8s.io/klog/v2 v2.30.0 // indirect
k8s.io/utils v0.0.0-20210930125809-cb0fa318a74b // indirect
sigs.k8s.io/structured-merge-diff/v4 v4.1.2 // indirect
sigs.k8s.io/yaml v1.2.0 // indirect
k8s.io/api v0.26.2 // indirect
k8s.io/apimachinery v0.26.2 // indirect
k8s.io/client-go v0.26.2 // indirect
k8s.io/klog/v2 v2.90.1 // indirect
k8s.io/kube-openapi v0.0.0-20230303024457-afdc3dddf62d // indirect
k8s.io/utils v0.0.0-20230220204549-a5ecb0141aa5 // indirect
sigs.k8s.io/json v0.0.0-20221116044647-bc3834ca7abd // indirect
sigs.k8s.io/structured-merge-diff/v4 v4.2.3 // indirect
sigs.k8s.io/yaml v1.3.0 // indirect
)
Loading

0 comments on commit ad5abb6

Please sign in to comment.