【问题标题】:MongoDB Kafka Source Connector throws java.lang.IllegalStateException: Queue full when using copy.existing: trueMongoDB Kafka 源连接器抛出 java.lang.IllegalStateException: Queue full when using copy.existing: true
【发布时间】:2020-05-29 08:01:09
【问题描述】:

使用连接器https://github.com/mongodb/mongo-kafka将数据从mongodb导入kafka时,会抛出java.lang.IllegalStateException: Queue full

我使用默认设置copy.existing.queue.size,即16000,和copy.existing: true。我应该设置什么值?集合大小为10G。

环境:

mongo-kafka-connect: 1.0.0
Kafka: 2.4.0
Kafka-Connect: 2.4.0
MongoDB server: 3.6.14
mongodb-driver-sync: 3.12.1

堆栈跟踪: org.apache.kafka.connect.errors.ConnectException: java.lang.IllegalStateException: Queue full\n\tat com.mongodb.kafka.connect.source.MongoCopyDataManager.poll(MongoCopyDataManager.java:95)\n\tat com.mongodb.kafka.connect.source.MongoSourceTask.getNextDocument(MongoSourceTask.java:301)\n\tat com.mongodb.kafka.connect.source.MongoSourceTask.poll(MongoSourceTask.java:154)\n\tat org.apache.kafka.connect.runtime.WorkerSourceTask.poll(WorkerSourceTask.java:265)\n\tat org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:232)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:177)\n\tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)\n\tat java.base\/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)\n\tat java.base\/java.util.concurrent.FutureTask.run(Unknown Source)\n\tat java.base\/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)\n\tat java.base\/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)\n\tat java.base\/java.lang.Thread.run(Unknown Source)\nCaused by: java.lang.IllegalStateException: Queue full\n\tat java.base\/java.util.AbstractQueue.add(Unknown Source)\n\tat java.base\/java.util.concurrent.ArrayBlockingQueue.add(Unknown Source)\n\tat com.mongodb.client.internal.Java8ForEachHelper.forEach(Java8ForEachHelper.java:30)\n\tat com.mongodb.client.internal.Java8AggregateIterableImpl.forEach(Java8AggregateIterableImpl.java:54)\n\tat com.mongodb.kafka.connect.source.MongoCopyDataManager.copyDataFrom(MongoCopyDataManager.java:123)\n\tat com.mongodb.kafka.connect.source.MongoCopyDataManager.lambda$new$0(MongoCopyDataManager.java:87)\n\t... 5 more

【问题讨论】:

    标签: apache-kafka mongodb-kafka-connector


    【解决方案1】:

    固定在 https://github.com/mongodb/mongo-kafka/commit/7e6bf97742f2ad75cde394d088823b86880cdf4e

    并将在 1.0.0 之后发布。所以如果有人遇到同样的问题,请将版本更新到 1.0.0 之后。

    【讨论】:

      猜你喜欢
      • 2020-04-27
      • 2023-01-01
      • 2021-08-11
      • 1970-01-01
      • 2019-12-27
      • 2019-07-22
      • 2020-05-17
      • 2021-11-23
      • 2021-04-06
      相关资源
      最近更新 更多