A MongoDB collection watcher that pushes oplog events into Kafka
50K+
This project listens for a MongoDB collection events (insert, update, delete, ...) also called "oplogs" for operation logs and distribute them into a Kafka topic of your choice. There is also a replay mode that allows you to initialize all items of a collection into a Kafka topic for the first time.
In addition of the binary, you will also need the following the Kafka library:
You can download the latest version of the binary built for your architecture here:
The watcher is also available as a Docker image. You can run it using the following example and pass configuration environment variables:
$ docker run \
-e 'KAFKA_MONGO_WATCHER_REPLAY=true' \
etf1/kafka-mongo-watcher:latest
Optionally, you can also download and build it from the sources. You have to retrieve the project sources by using one of the following way:
$ go get -u github.com/etf1/kafka-mongo-watcher
# or
$ git clone https://github.com/etf1/kafka-mongo-watcher.git
Then, build the binary:
$ GOOS=linux GOARCH=amd64 go build -ldflags '-s -w' -o kafka-mongo-watcher ./cmd/watcher/
In order to run the watcher, type the following command with the desired arguments.
You can use flags (as in this example) or environment variables:
$ ./kafka-mongo-watcher -KAFKA_MONGO_WATCHER_REPLAY=true
...
<info> Tech HTTP server started {"facility":"kafka-mongo-watcher","version":"wip","addr":":8001","file":"/usr/local/Cellar/go/1.14/libexec/src/runtime/asm_amd64.s","line":1373}
<info> Connected to mongodb database {"facility":"kafka-mongo-watcher","version":"wip","uri":"mongodb://root:[email protected]:27011,127.0.0.1:27012,127.0.0.1:27013/watcher?replicaSet=replicaset\u0026authSource=admin"}
<info> Connected to kafka producer {"facility":"kafka-mongo-watcher","version":"wip","bootstrao-servers":"127.0.0.1:9092"}
...
In dev environment you can copy .env.dist in .env and edit his content in order to customize easily the env variables.
You can set/override configuration variables from .env file and from variables environment and or from cli arguments
(If a variables was configured in multiple sources the last will override the previous one)
Type: string
Description: In case you want to specify a filtering pipeline, you can specify it here. It works both wil replay and watch mode.
Example value: [ { "$match": { "fullDocument.is_active": true } }, { $addFields: { "custom-field": "custom-value" } } ]
Type: bool
Description: In case you want to send all collection's documents once (default: false)
Type: string
Description: The MongoDB connection string URI (default: mongodb://root:[email protected]:27011,...)
Type: string
Description: The MongoDB collection you want to watch (default: "items")
Type: string
Description: The MongoDB database name you want to connect to (default: "watcher")
Type: duration
Description: The MongoDB server selection timeout duration (default: 2s)
Type: integer
Description: In case you want to enable watch batch size on MongoDB watch (default: 0 / no batch)
Type: boolean
Description: In case you want to retrieve the full document when watching for oplogs (default: true)
Type: duration
Description: In case you want to set a maximum value awaiting for new oplogs (default: 0 / don't stop)
Type: string
Description: In case you want to set a logical starting point for the change stream (example : {"_data": <hex string>})
Type: uint32 (increment value)
Type: uint32 (timestamp)
Description: In case you want to set a timestamp for the change stream to only return changes that occurred at or after the given timestamp (default: nil)
Type: integer
Description: The max number of retries when trying to watch a collection (default: 3, set to 0 to disable retry)
Type: duration
Description: Sleeping delay between two watch attempts (default: 500ms)
Type: string
Description: Kafka bootstrap servers list (default: "127.0.0.1:9092")
Type: string
Description: Kafka topic to write into (default: "kafka-mongo-watcher")
Type: integer
Description: The maximum size of the internal channel producer size (default: 10000)
A big value here can increase the heap memory of the application as all the payload that have to be sent to Kafka will be maintained in channel.
Type: boolean
Description: Used to enable/disable log verbosity (default: true)
Type: string
Description: Used to define first level you want to start display logs (default: "info")
Type: string
Description: In case you want to push logs into a Graylog server, just fill this entry with the endpoint
Type: duration
Description: A idle timeout for HTTP technical server (default: 90s)
Type: duration
Description: A read timeout for HTTP technical server (default: 1s)
Type: duration
Description: A write timeout for HTTP technical server (default: 10s)
Type: string
Description: A specified address for HTTP technical server to listen (default: ":8001")
Type: boolean
Description: Used to enable/disable the configuration print at startup (default: true)
Type: boolean
Description: In case you want to enable Go pprof debugging (default: true). No impact when not used
The watcher also exposes metrics about Go process and Watcher application.
These metrics can be scraped by Prometheus by browsing the following technical HTTP server endpoint: http://127.0.0.1:8001/metrics
Unit tests can be run with the following command:
$ go test -v -mod vendor ./...
And integration tests can be run with:
$ make test-integration
This will load needed mongodb and kafka containers and run the tests suite
Content type
Image
Digest
sha256:b086d545c…
Size
15.7 MB
Last updated
21 days ago
docker pull etf1/kafka-mongo-watcher:v0.7.1-alpha