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

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及Oracle ojdbc17.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配置错误,还是主题授权问题?


排查思路

  1. SASL认证配置缺失:KafkaConnect的全局认证配置不会自动传递给Debezium的历史主题生产者/消费者。当前Connector仅配置了consumer/producer.sasl.mechanism和security.protocol,但缺少consumer/producer.sasl.jaas.config,导致历史主题客户端无法完成SASL认证,进而无法连接Kafka获取元数据。
  2. 主题创建权限不足:即使Kafka启用自动创建主题,admin用户可能没有Create权限。需检查Strimzi的KafkaUser资源权限配置,确认是否包含Write和Create权限。
  3. 历史主题客户端配置不完整:Debezium的database.history和schema.history.internal组件使用独立Kafka客户端,需完整安全配置。要确保bootstrap.servers的协议与Kafka监听配置一致(如SASL_PLAINTEXT)。
  4. Avro转换器间接影响:若Schema Registry连接异常可能导致Connector初始化失败,但日志无相关报错,优先级低于认证问题,可先排除认证问题后再验证。

解决方案

  1. 补充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
  1. 预创建历史主题:若自动创建仍有问题,可手动创建schema-changes和history-changes主题,设置正确分区数和副本数,确认Connector能正常写入。
  2. 检查bootstrap地址协议:确保database.history.kafka.bootstrap.servers和schema.history.internal.kafka.bootstrap.servers使用的协议与Kafka监听配置完全匹配,避免协议不兼容导致连接失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 23:18:14