Kafka Connect in distributed mode on Kubernetes received unknown topic or partition error in fetch for partition

Viewed 22

I'm trying to deploy Debezium source connector on Kafka Connect 3.1.0 in distributed mode. I deployed 3 Kafka brokers with 3 Zookeeper instances and built Kafka Connect pod from Debezium tutorial. I dived into Kafka Connect pod logs and encountered those warnings:

2022-09-09 12:23:22,392 WARN [Consumer clientId=consumer-connect-cluster-1, groupId=connect-cluster] Received unknown topic or partition error in fetch for partition connect-cluster-offsets-10 (org.apache.kafka.clients.consumer.internals.Fetcher) [KafkaBasedLog Work Thread - connect-cluster-offsets]

while logs from running connect-distributed.sh show this:

[2022-09-09 12:39:29,570] INFO [Worker clientId=connect-1, groupId=connect-cluster] Rebalance started (org.apache.kafka.connect.runtime.distributed.WorkerCoordinator:222)
[2022-09-09 12:39:29,570] INFO [Worker clientId=connect-1, groupId=connect-cluster] (Re-)joining group (org.apache.kafka.connect.runtime.distributed.WorkerCoordinator:535)
[2022-09-09 12:39:29,571] INFO [Worker clientId=connect-1, groupId=connect-cluster] Successfully joined group with generation Generation{generationId=13803870, memberId='connect-1-fb18f92c-6c73-461b-8b42-4bdee3e91e99', protocol='sessioned'} (org.apache.kafka.connect.runtime.distributed.WorkerCoordinator:595)
[2022-09-09 12:39:29,575] INFO [Worker clientId=connect-1, groupId=connect-cluster] Successfully synced group in generation Generation{generationId=13803870, memberId='connect-1-fb18f92c-6c73-461b-8b42-4bdee3e91e99', protocol='sessioned'} (org.apache.kafka.connect.runtime.distributed.WorkerCoordinator:761)
[2022-09-09 12:39:29,575] INFO [Worker clientId=connect-1, groupId=connect-cluster] Joined group at generation 13803870 with protocol version 2 and got assignment: Assignment{error=0, leader='connect-1-91c7dc97-e954-4b7e-b3a3-25021ba65dc3', leaderUrl='http://10.244.234.90:8083/', offset=1, connectorIds=[], taskIds=[], revokedConnectorIds=[], revokedTaskIds=[], delay=0} with rebalance delay: 0 (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1853)
[2022-09-09 12:39:29,575] INFO [Worker clientId=connect-1, groupId=connect-cluster] Starting connectors and tasks using config offset 1 (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1378)
[2022-09-09 12:39:29,575] INFO [Worker clientId=connect-1, groupId=connect-cluster] Finished starting connectors and tasks (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1406)

I' ve changed connect-distributed.properties multiple times, as well as modifying curl command to deploy a connector. I especially experimented with listeners and advertised.listeners fields by changing service names, ports, protocols. So far, the connect-distributed.properties file looks like this (meaning connect-distributed.sh doesn't crash):

bootstrap.servers=kafka-whatever-kafka-bootstrap:9092
offset.flush.interval.ms=10000
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true
offset.storage.topic=connect-cluster-offsets
offset.storage.replication.factor=1
config.storage.topic=connect-cluster-configs
config.storage.replication.factor=1
status.storage.topic=connect-cluster-status
status.storage.replication.factor=1
plugin.path=/opt/kafka/plugins
listeners=http://0.0.0.0:9092,http://0.0.0.0:9090,http://0.0.0.0:9091,http://0.0.0.0:9093
advertised.listeners=http://connect-cluster-connect-api:8083
group.id=connect-cluster

curl command to deploy a connector that I used is:

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://127.0.0.1:8080/connectors/ -d '{ "name": "inventory-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "184054", "database.server.name": "dbserver1", "database.include.list": "inventory", "database.history.kafka.bootstrap.servers": "kafka-whatever-kafka-bootstrap:9092", "database.history.kafka.topic": "schema-changes.inventory" } }'
0 Answers
Related