You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

MongoDB到Event Hub的Kafka Connect连接器初始失败后无法重连发消息

Kafka Connect对接Azure Event Hub大消息发送失败问题

我们在Kubernetes环境中通过Kafka Connect将MongoDB源连接器对接至Azure Event Hub,测试大消息场景时,连接器成功发送若干条消息后出现失败;后续尝试为Kafka Connect Pod扩容资源,或通过Kafka Connect端点删除并重建连接器,问题均未解决。

错误日志

15:27:05,982 INFO
[cosmosdb-source-develop|task-0] Monitor thread successfully connected to server with description
•   ServerDescription:
o   address: <my_server_address>
o   type: REPLICA_SET_PRIMARY
o   state: CONNECTED
o   ok: true
o   minWireVersion: 0
o   maxWireVersion: 8
o   maxDocumentSize: 16777216
o   logicalSessionTimeoutMinutes: 30
o   roundTripTimeNanos: 86778707
o   setName: '<my_set_name>'
o   canonicalAddress: <my_server_address>
o   hosts: [<my_server_address>]
o   tagSet: TagSet{[Tag{name='region', value='<my_region>' }]}
o   setVersion: 1
o   lastUpdateTimeNanos: 12352499337633898
(org.mongodb.driver.cluster)
[cluster-ClusterId{value='<my_cluster_id>'}-<my_server_address>]

15:27:05,982 INFO
[cosmosdb-source-develop|task-0] Discovered replica set primary <my_server_address>
(org.mongodb.driver.cluster)
[cluster-ClusterId{value='<my_cluster_id>'}-<my_server_address>]

15:27:05,993 INFO
[cosmosdb-source-develop|task-0] Opened connection [connectionId{localValue:4, serverValue:863417751}] to <my_server_address>
(org.mongodb.driver.connection)
[cluster-rtt-ClusterId{value='<my_cluster_id>'}-<my_server_address>]

15:27:06,154 INFO
[cosmosdb-source-develop|task-0] Opened connection [connectionId{localValue:5, serverValue:366833602}] to <my_server_address>
(org.mongodb.driver.connection)
[task-thread-cosmosdb-source-develop-0]

15:27:06,168 INFO
[cosmosdb-source-develop|task-0] Watching for collection changes on '<my_db>.<my_collection>'
(com.mongodb.kafka.connect.source.MongoSourceTask)
[task-thread-cosmosdb-source-develop-0]

15:27:06,172 INFO
[cosmosdb-source-develop|task-0] Resuming the change stream after the previous offset using resumeAfter:
•   resumeAfter: {"_data": {"$binary": {"base64": "eyJWIjoyLCJSaWQiOiJnd0pQQUxGTGRuUT0iLCJDb250aW51YXRpb24iOlt7IkZlZWRSYW5nZSI6eyJ0eXBlIjoiRWZmZWN0aXZlIFBhcnRpdGlvbiBLZXkgUmFuZ2UiLCJ2YWx1ZSI6eyJtaW4iOiIiLCJtYXgiOiJGRiJ9fSwiU3RhdGUiOnsidHlwZSI6ImNvbnRpbnVhdGlvbiIsInZhbHVlIjoiXCIzMzdcIiJ9fV19", "subType": "00"}}, "_kind": 1}
(com.mongodb.kafka.connect.source.MongoSourceTask)
[task-thread-cosmosdb-source-develop-0]

15:27:06,752 INFO
[cosmosdb-source-develop|task-0] Started MongoDB source task
(com.mongodb.kafka.connect.source.MongoSourceTask)
[task-thread-cosmosdb-source-develop-0]

15:27:06,752 INFO
[cosmosdb-source-develop|task-0] WorkerSourceTask{id=cosmosdb-source-develop-0} Source task finished initialization and start
(org.apache.kafka.connect.runtime.AbstractWorkerSourceTask)
[task-thread-cosmosdb-source-develop-0]
15:27:12,864 INFO
[cosmosdb-source-develop|task-0] [Producer clientId=connector-producer-cosmosdb-source-develop-0] Node 0 disconnected.
(org.apache.kafka.clients.NetworkClient)
[kafka-producer-network-thread | connector-producer-cosmosdb-source-develop-0]

15:27:12,868 INFO
[cosmosdb-source-develop|task-0] [Producer clientId=connector-producer-cosmosdb-source-develop-0] Cancelled in-flight PRODUCE request with correlation id 7 due to node 0 being disconnected (elapsed time since creation: 8ms, elapsed time since send: 8ms, request timeout: 30000ms)
(org.apache.kafka.clients.NetworkClient)
[kafka-producer-network-thread | connector-producer-cosmosdb-source-develop-0]

15:27:12,869 WARN
[cosmosdb-source-develop|task-0] [Producer clientId=connector-producer-cosmosdb-source-develop-0] Got error produce response with correlation id 7 on topic-partition develop1.<my_db>.<my_collection>-16, retrying (2147483646 attempts left). Error: NETWORK_EXCEPTION. Error Message: Disconnected from node 0
(org.apache.kafka.clients.producer.internals.Sender)
[kafka-producer-network-thread | connector-producer-cosmosdb-source-develop-0]

15:27:12,870 WARN
[cosmosdb-source-develop|task-0] [Producer clientId=connector-producer-cosmosdb-source-develop-0] Received invalid metadata error in produce request on partition develop1.<my_db>.<my_collection>-16 due to org.apache.kafka.common.errors.NetworkException: Disconnected from node 0. Going to request metadata update now
(org.apache.kafka.clients.producer.internals.Sender)
[kafka-producer-network-thread | connector-producer-cosmosdb-source-develop-0]

配置信息

ProducerConfig配置

2024-10-21 15:27:05,277 INFO [cosmosdb-source-develop|task-0] ProducerConfig values:
        acks = -1
        auto.include.jmx.reporter = true
        batch.size = 16384
        bootstrap.servers = <my bootstrap server> 
        buffer.memory = 33554432
        client.dns.lookup = use_all_dns_ips
        client.id = connector-producer-cosmosdb-source-develop-0
        compression.type = none
        connections.max.idle.ms = 540000
        delivery.timeout.ms = 120000
        enable.idempotence = true
        enable.metrics.push = true
        interceptor.classes = []
        key.serializer = class org.apache.kafka.common.serialization.ByteArraySerializer
        linger.ms = 0
        max.block.ms = 9223372036854775807
        max.in.flight.requests.per.connection = 1
        max.request.size = 80518918
        metadata.max.age.ms = 300000
        metadata.max.idle.ms = 300000
        metric.reporters = []
        metrics.num.samples = 2
        metrics.recording.level = INFO
        metrics.sample.window.ms = 30000
        partitioner.adaptive.partitioning.enable = true
        partitioner.availability.timeout.ms = 0
        partitioner.class = null
        partitioner.ignore.keys = false
        receive.buffer.bytes = 32768
        reconnect.backoff.max.ms = 1000
        reconnect.backoff.ms = 50
        request.timeout.ms = 30000
        retries = 2147483647
        retry.backoff.max.ms = 1000
        retry.backoff.ms = 100
        sasl.client.callback.handler.class = null
        sasl.jaas.config = [hidden]
        sasl.kerberos.kinit.cmd = /usr/bin/kinit
        sasl.kerberos.min.time.before.relogin = 60000
        sasl.kerberos.service.name = null
        sasl.kerberos.ticket.renew.jitter = 0.05
        sasl.kerberos.ticket.renew.window.factor = 0.8
        sasl.login.callback.handler.class = null
        sasl.login.class = null
        sasl.login.connect.timeout.ms = null
        sasl.login.read.timeout.ms = null
        sasl.login.refresh.buffer.seconds = 300
        sasl.login.refresh.min.period.seconds = 60
        sasl.login.refresh.window.factor = 0.8
        sasl.login.refresh.window.jitter = 0.05
        sasl.login.retry.backoff.max.ms = 10000
        sasl.login.retry.backoff.ms = 100
        sasl.mechanism = PLAIN
        sasl.oauthbearer.clock.skew.seconds = 30
        sasl.oauthbearer.expected.audience = null
        sasl.oauthbearer.expected.issuer = null
        sasl.oauthbearer.jwks.endpoint.refresh.ms = 3600000
        sasl.oauthbearer.jwks.endpoint.retry.backoff.max.ms = 10000
        sasl.oauthbearer.jwks.endpoint.retry.backoff.ms = 100
        sasl.oauthbearer.jwks.endpoint.url = null
        sasl.oauthbearer.scope.claim.name = scope
        sasl.oauthbearer.sub.claim.name = sub
        sasl.oauthbearer.token.endpoint.url = null
        security.protocol = <my_security_protol>
        security.providers = null
        send.buffer.bytes = 131072
        socket.connection.setup.timeout.max.ms = 30000
        socket.connection.setup.timeout.ms = 10000
        <my_ssl_details> 
        transaction.timeout.ms = 60000
        transactional.id = null
        value.serializer = class org.apache.kafka.common.serialization.ByteArraySerializer

SourceConnectorConfig配置

INFO SourceConnectorConfig values:
        config.action.reload = restart
        connector.class = com.mongodb.kafka.connect.MongoSourceConnector
        errors.log.enable = false
        errors.log.include.messages = false
        errors.retry.delay.max.ms = 60000
        errors.retry.timeout = 0
        errors.tolerance = none
        exactly.once.support = requested
        header.converter = null
        key.converter = null
        name = cosmosdb-source-develop
        offsets.storage.topic = null
        predicates = []
        tasks.max = 1
        topic.creation.groups = []
        transaction.boundary = poll
        transaction.boundary.interval.ms = null
        transforms = []
        value.converter = null
 
(org.apache.kafka.connect.runtime.SourceConnectorConfig) [DistributedHerder-<info>]
2024-10-21 15:27:05,356 INFO EnrichedConnectorConfig values:
        config.action.reload = restart
        connector.class = com.mongodb.kafka.connect.MongoSourceConnector
        errors.log.enable = false
        errors.log.include.messages = false
        errors.retry.delay.max.ms = 60000
        errors.retry.timeout = 0
        errors.tolerance = none
        exactly.once.support = requested
        header.converter = null
        key.converter = null
        name = cosmosdb-source-develop
        offsets.storage.topic = null
        predicates = []
        tasks.max = 1
        topic.creation.groups = []
        transaction.boundary = poll
        transaction.boundary.interval.ms = null
        transforms = []
        value.converter = null

内容的提问来源于stack exchange,提问作者gt1

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 00:42:01