Strimzi+Debezium环境下KafkaConnector无法自动创建主题问题
KafkaConnector无法自动创建主题问题排查
环境与技术栈
- Kafka:基于Strimzi Operator
0.46.0,KRaft模式,版本4.0.0 - KafkaConnect:使用
quay.io/strimzi/kafka:0.46.0-kafka-4.0.0镜像,扩展了Debezium JAR、Avro转换器JAR及Oracleojdbc17.jar - Schema Registry:
docker.io/bitnami/schema-registry:7.9.0 - Kafka已启用
auto.create.topics.enable: true
问题现象
KafkaConnect集群可正常创建内部主题(connect-cluster-configs、connect-cluster-offsets、connect-cluster-status),但部署Debezium Oracle源连接器后,无法自动创建schema.history.internal.kafka.topic、database.history.kafka.topic及数据表对应的主题,日志报元数据获取超时:
2025-06-03 14:56:28 INFO [kafka-producer-network-thread | next-schemahistory] NetworkClient:411 - [Producer clientId=next-schemahistory] Cancelled in-flight METADATA request with correlation id 126 due to node -1 being disconnected (elapsed time since creation: 1ms, elapsed time since send: 1ms, throttle time: 0ms, request timeout: 30000ms) 2025-06-03 14:56:28 WARN [kafka-producer-network-thread | next-schemahistory] NetworkClient:1255 - [Producer clientId=next-schemahistory] Bootstrap broker kafka-kafka-bootstrap:9092 (id: -1 rack: null isFenced: false) disconnected 2025-06-03 14:56:28 INFO [SourceTaskOffsetCommitter-1] BaseSourceTask:503 - Couldn't commit processed log positions with the source database due to a concurrent connector shutdown or restart 2025-06-03 14:56:28 INFO [task-thread-debezium-connector-ko-employee-0] ConsumerCoordinator:1056 - [Consumer clientId=next-schemahistory, groupId=next-schemahistory] Resetting generation and member id due to: consumer pro-actively leaving the group 2025-06-03 14:56:28 INFO [task-thread-debezium-connector-ko-employee-0] ConsumerCoordinator:1103 - [Consumer clientId=next-schemahistory, groupId=next-schemahistory] Request joining group due to: consumer pro-actively leaving the group 2025-06-03 14:56:28 INFO [task-thread-debezium-connector-ko-employee-0] AppInfoParser:89 - App info kafka.consumer for next-schemahistory unregistered 2025-06-03 14:56:28 ERROR [task-thread-debezium-connector-ko-employee-0] WorkerTask:234 - WorkerSourceTask{id=debezium-connector-ko-employee-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata 2025-06-03 14:56:28 INFO [task-thread-debezium-connector-ko-employee-0] BaseSourceTask:436 - Stopping down connector 2025-06-03 14:56:28 INFO [pool-14-thread-1] JdbcConnection:983 - Connection gracefully closed 2025-06-03 14:56:28 INFO [pool-15-thread-1] JdbcConnection:983 - Connection gracefully closed 2025-06-03 14:56:28 INFO [task-thread-debezium-connector-ko-employee-0] KafkaProducer:1367 - [Producer clientId=next-schemahistory] Closing the Kafka producer with timeoutMillis = 30000 ms.
最终连接器自动关闭,无目标主题生成。
配置详情
所有服务使用同一SASL用户admin,认证配置如下:
security.protocol: SASL_PLAINTEXT sasl.mechanism: SCRAM-SHA-512 sasl.jaas.config: ${secrets:admin/sasl.jaas.config}
KafkaConnector配置:
--- apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnector metadata: name: debezium-connector-ko-employee namespace: {{ .Release.Namespace }} labels: strimzi.io/cluster: {{ .Values.debezium.cluster.name }} annotations: strimzi.io/use-connector-resources: "true" argocd.argoproj.io/sync-wave: "6" spec: class: io.debezium.connector.oracle.OracleConnector tasksMax: 1 autoRestart: enabled: true config: oracle.connection.string: "jdbc:oracle:thin:@//<WORKING_JDBC_STRING>" database.hostname: "<HOST>" database.url: "jdbc:oracle:thin:@//<WORKING_JDBC_STRING>" database.port: "<PORT>" database.user: "<USER>" database.password: "<PASS>" # ... Think about using a Secret database.dbname: "<DB_NAME>" database.pdb.name: "<PDB_NAME>" table.include.list: "SFNKO.MITARBEITER_OUTBOX,SFNKO.MITARBEITER" schema.include.list: "SFNKO" database.include.list: "<DBs>" database.history.kafka.bootstrap.servers: {{ .Values.kafka.authentication.security_protocol }}://{{ .Release.Name }}-kafka-bootstrap:{{ .Values.kafka.ports.plain }} database.history.kafka.topic: "schema-changes" schema.history.internal.kafka.bootstrap.servers: {{ .Values.kafka.authentication.security_protocol }}://{{ .Release.Name }}-kafka-bootstrap:{{ .Values.kafka.ports.plain }} schema.history.internal.kafka.topic: "history-changes" database.history.store.only.captured.tables.ddl: "true" include.schema.changes: "false" topic.prefix: "next" snapshot.mode: "when_needed" log.mining.strategy: "online_catalog" schema.history.internal.store.only.captured.tables.ddl: "true" schema.history.internal.store.only.captured.databases.ddl: "true" errors.log.include.messages: "true" cdc.flattening.enabled: "true" key.converter: "io.confluent.connect.avro.AvroConverter" key.converter.schema.registry.url: "http://{{ .Values.schema_registry.name }}:{{ .Values.schema_registry.port }}" key.converter.schema.registry.auto-register: "true" key.converter.schema.registry.find-latest: "true" value.converter: "io.confluent.connect.avro.AvroConverter" value.converter.schema.registry.url: "http://{{ .Values.schema_registry.name }}:{{ .Values.schema_registry.port }}" value.converter.schema.registry.auto-register: "true" value.converter.schema.registry.find-latest: "true" schema.name.adjustment.mode: "avro" signal.data.collection: "<PREFIX>.SFNLNK.DEBEZIUM_SIGNAL" transforms: "changes,unwrap" transforms.changes.type: "io.debezium.transforms.ExtractChangedRecordState" transforms.changes.header.changed.name: "Changed" transforms.changes.header.unchanged.name: "Unchanged" transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState" transforms.unwrap.drop.tombstones: "true" transforms.unwrap.delete.handling.mode: "rewrite" transforms.unwrap.add.fields: "op" incremental.snapshot.chunk.size: "262144" max.batch.size: "16384" max.queue.size: "65536" snapshot.max.threads: "1" topic.creation.default.replication.factor: "1" topic.creation.default.partitions: "1" topic.creation.default.cleanup.policy: "compact" topic.creation.default.compression.type: "lz4" consumer.sasl.mechanism: {{ upper .Values.kafka.authentication.type }} consumer.security.protocol: {{ .Values.kafka.authentication.security_protocol }} producer.sasl.mechanism: {{ upper .Values.kafka.authentication.type }} producer.security.protocol: {{ .Values.kafka.authentication.security_protocol }}
已验证内容
- KafkaConnect使用相同凭据成功创建内部主题
- KafkaConnector启动后因元数据超时失败
- 集群内可正常访问Kafka及Schema Registry
- 连接器Pod可手动连通Kafka bootstrap地址
核心疑问
为何KafkaConnect能正常创建内部主题,但KafkaConnector无法自动生成schema-changes、history-changes及数据表对应主题?是Avro转换器配置、网络/SASL配置错误,还是主题授权问题?
排查思路
- SASL认证配置缺失:KafkaConnect的全局认证配置不会自动传递给Debezium的历史主题生产者/消费者。当前Connector仅配置了
consumer/producer.sasl.mechanism和security.protocol,但缺少consumer/producer.sasl.jaas.config,导致历史主题客户端无法完成SASL认证,进而无法连接Kafka获取元数据。 - 主题创建权限不足:即使Kafka启用自动创建主题,
admin用户可能没有Create权限。需检查Strimzi的KafkaUser资源权限配置,确认是否包含Write和Create权限。 - 历史主题客户端配置不完整:Debezium的
database.history和schema.history.internal组件使用独立Kafka客户端,需完整安全配置。要确保bootstrap.servers的协议与Kafka监听配置一致(如SASL_PLAINTEXT)。 - Avro转换器间接影响:若Schema Registry连接异常可能导致Connector初始化失败,但日志无相关报错,优先级低于认证问题,可先排除认证问题后再验证。
解决方案
- 补充Connector的SASL JAAS配置:在KafkaConnector的
spec.config中添加以下参数:
consumer.sasl.jaas.config: ${secrets:admin/sasl.jaas.config} producer.sasl.jaas.config: ${secrets:admin/sasl.jaas.config}
若使用Strimzi的use-connector-resources注解,也可将JAAS配置挂载到Connect Pod环境变量中,确保Debezium客户端能读取到。
2. 验证Kafka用户权限:检查admin用户的KafkaUser配置,确保包含以下权限:
authorization: type: simple acls: - resource: type: topic name: "*" patternType: literal operation: Create - resource: type: topic name: "*" patternType: literal operation: Write - resource: type: topic name: "*" patternType: literal operation: Read
- 预创建历史主题:若自动创建仍有问题,可手动创建
schema-changes和history-changes主题,设置正确分区数和副本数,确认Connector能正常写入。 - 检查bootstrap地址协议:确保
database.history.kafka.bootstrap.servers和schema.history.internal.kafka.bootstrap.servers使用的协议与Kafka监听配置完全匹配,避免协议不兼容导致连接失败。
内容的提问来源于stack exchange,提问作者Sebastian Sommerfeld
相关产品推荐
相关产品推荐

