Mongocxx change_stream callback

Viewed 167

I have a process that is monitoring a mongo db and needs to be notified when there is a change to a table.

It seems like the most logical way to handle this is via a change_stream with a callback that executes when something has changed - that is, the same way a "watch" functions in JavaScript.

Looking at the documentation, I see that the changestream object has iterators, but I don't see the option for a callback. Is there a way to handle this and/or are there any good examples to leverage ?

1 Answers

Mongocxx does not have a callback API, but since the change stream API has blocking calls, you can implement your own callback system with a class that spawns a thread for watching, and then invokes callbacks when something happens. Example code below:

#include <mongocxx/client.hpp>
#include <mongocxx/instance.hpp>
#include <bsoncxx/document/view.hpp>

// STD
#include <iostream>
#include <string>
#include <thread>
#include <functional>
#include <chrono>

class MongoWatcher 
{
    public:
        MongoWatcher(std::function<void(std::string)> callback);

    private:
        std::thread* m_changesThread;
        std::function<void(std::string)> m_changeCallback;
 
        void watchCollection();
};

MongoWatcher::MongoWatcher(std::function<void(std::string)> callback)
{
    // Asssuming Mongocxx instance has been created by invoking code
    m_changeCallback = callback;
    m_changesThread = new std::thread([this](){watchCollection();});
}

void MongoWatcher::watchCollection()
{
    mongocxx::uri uri("mongodb://localhost:27017");
    mongocxx::client client = mongocxx::client(uri);
    mongocxx::database db = client["db"];
    mongocxx::collection collection = db["col"];
    while (true) // Loop forever
    {
        try
        {
            mongocxx::options::change_stream options;
            options.full_document("updateLookup");
            const std::chrono::milliseconds await_time{1000};
            options.max_await_time(await_time);
            mongocxx::change_stream streamChanges = collection.watch(options);
            for (const bsoncxx::document::view& event : streamChanges)
            {
                bsoncxx::document::view view = event;
                bsoncxx::document::element element = view["operationType"];
                bsoncxx::types::b_utf8 utf = element.get_utf8();
                std::string operationType = utf.value.data();
                m_changeCallback(operationType);
            }
        }
        catch (const std::exception& e)
        {
            std::cerr << "MongoDB watcher caught exception: " << e.what() << std::endl;
        }
        // Take a short pause before checking for changes again
        std::this_thread::sleep_for(std::chrono::milliseconds(500));
    }
}
Related