In this article, we will see how to solve org.apache.kafka.connect.errors.ConnectException: An exception occurred when trying to get the next item from the change stream: Query failed with error code 17287 with name ' ' and error message 'Bad query specified'. If you are also getting this error then there is a chance that you are also using kafka connector of different database than the one you are connecting to. This sometimes may result in incompatibility which causes this error. The incompatibility gets triggered due to many reasons. One reason could be the patching of database cluster engine which was exactly the case with me. Let's understand this in detail.

Solved "An exception occurred when trying to get the next item from the change stream: Query failed with error code 17287"
Also Read: How to republish existing data using Apache Kafka Connect
So as I was explaining, I was having a setup where mongodb kafka connector connecting to documentdb. While the connector was running fine, I noticed that task running on connector started failing.
cyberithub@macos1066 % kcctl get connectors NAME TYPE STATE TASKS cyberithub-connector source RUNNING 0: FAILED
When I tried to deep dive into error, I noticed below error in description.
cyberithub@macos1066 % kcctl describe connector cyberithub-connector Name: cyberithub-connector Type: source State: RUNNING Worker ID: example.cyberithub.com:8083 Config: change.stream.full.document: updateLookup collection: testcol connection.uri: mongodb://${file:/usr/share/cyberithub/secrets.properties:username}:${file:/usr/share/cyberithub/secrets.properties:password}@cyberithub.example.server.com:27017/?ssl=true connector.class: com.mongodb.kafka.connect.MongoSourceConnector database: testdb errors.tolerance: none name: cyberithub-connector producer.override.max.request.size: 2700000 tasks.max: 1 topic.prefix: cyber.col Tasks: 0: State: FAILED Worker ID: example.cyberithub.com:8083 Trace: org.apache.kafka.connect.errors.ConnectException: An exception occurred when trying to get the next item from the Change Stream: Query failed with error code 17287 with name '' and error message 'Bad query specified' on server example.cyberithub.com:8083 at com.mongodb.kafka.connect.source.StartedMongoSourceTask.getNextBatch(StartedMongoSourceTask.java:612) at com.mongodb.kafka.connect.source.StartedMongoSourceTask.pollInternal(StartedMongoSourceTask.java:213) at com.mongodb.kafka.connect.source.StartedMongoSourceTask.poll(StartedMongoSourceTask.java:190) at com.mongodb.kafka.connect.source.MongoSourceTask.poll(MongoSourceTask.java:180) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.poll(AbstractWorkerSourceTask.java:469) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:357) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259) at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:77) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:237) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833) Caused by: com.mongodb.MongoQueryException: Query failed with error code 17287 with name '' and error message 'Bad query specified' on server example.cyberithub.com:8083 at com.mongodb.internal.operation.QueryHelper.translateCommandException(QueryHelper.java:29) at com.mongodb.internal.operation.QueryBatchCursor.lambda$getMore$1(QueryBatchCursor.java:298) at com.mongodb.internal.operation.QueryBatchCursor$ResourceManager.executeWithConnection(QueryBatchCursor.java:532) at com.mongodb.internal.operation.QueryBatchCursor.getMore(QueryBatchCursor.java:286) at com.mongodb.internal.operation.QueryBatchCursor.tryHasNext(QueryBatchCursor.java:239) at com.mongodb.internal.operation.QueryBatchCursor.lambda$tryNext$0(QueryBatchCursor.java:222) at com.mongodb.internal.operation.QueryBatchCursor$ResourceManager.execute(QueryBatchCursor.java:417) at com.mongodb.internal.operation.QueryBatchCursor.tryNext(QueryBatchCursor.java:221) at com.mongodb.internal.operation.ChangeStreamBatchCursor$3.apply(ChangeStreamBatchCursor.java:107) at com.mongodb.internal.operation.ChangeStreamBatchCursor$3.apply(ChangeStreamBatchCursor.java:103) at com.mongodb.internal.operation.ChangeStreamBatchCursor.resumeableOperation(ChangeStreamBatchCursor.java:200) at com.mongodb.internal.operation.ChangeStreamBatchCursor.tryNext(ChangeStreamBatchCursor.java:103) at com.mongodb.client.internal.MongoChangeStreamCursorImpl.tryNext(MongoChangeStreamCursorImpl.java:87) at com.mongodb.kafka.connect.source.StartedMongoSourceTask.getNextBatch(StartedMongoSourceTask.java:594) ... 14 more Topics: cyber.col.topic
After checking for a while, I noticed in database event that documentdb database cluster engine patch version was upgraded. This was found to be the key reason for task failure on kafka connector cyberithub-connector. If you are also facing this kind of situation then all you have to do is to just stop the connector and then resume it back to start processing from same offset without losing any event. To stop the connector, run kcctl stop connector cyberithub-connector as shown below.
cyberithub@macos1066 % kcctl stop connector cyberithub-connector Stopped connector cyberithub-connector
Then resume it by using kcctl resume connector cyberithub-connector command as shown below.
cyberithub@macos1066 % kcctl resume connector cyberithub-connector Resumed connector cyberithub-connector
Once done, if you check the status again using kcctl get connectors command. You should see task coming back to running state.
cyberithub@macos1066 % kcctl get connectors NAME TYPE STATE TASKS cyberithub-connector source RUNNING 0: RUNNING
This confirms that org.apache.kafka.connect.errors.ConnectException: An exception occurred when trying to get the next item from the change stream: Query failed with error code 17287 with name ' ' and error message 'Bad query specified' error is resolved now. It is however important to mention here that some of you might be thinking why we are not simply restarting the connector instead of stopping and resuming the connector? Well, you could do that but the thing is restart does not work here atleast not in my case. I tried that too but nothing happened, error remains the same.
Second thing is stopping and resuming connector is much more safer than restarting the connector especially in critical environments where loss of events can be catastrophic sometimes. When we restart a connector, task usually gets recreated but when we stop a connector, task gets paused and connector remains configured properly in kafka connect. So that when we resume again, it would start processing from the same state. Having said that, sometimes restarting connector is also required especially when we are reconfiguring the connector and would like to reinitialize the task. Hope this makes sense. Please let me know in comment box if this helps you.