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

Strimzi部署KafkaConnect对接Azure Event Hub时Schema Registry报错

问题描述

我在Kubernetes中通过Strimzi部署Kafka Connect,用Azure Event Hub替代原生Kafka作为Broker,尝试将发送到Event Hub的负载序列化为Avro格式。容器部署后插件加载正常,但Schema Registry出现报错。

已验证:不启用Avro序列化时,能正常向Event Hub发送消息,问题明确出在Schema Registry配置上。

当前KafkaConnect配置如下:

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaConnect
metadata:
  name: my-connect-cluster
  annotations:
    strimzi.io/use-connector-resources: "false"
spec:
  version: 3.7.1
  replicas: 1
  bootstrapServers: mynamespace.servicebus.windows.net:9093 
  authentication:
    type: plain 
    username: $MyConnectionString
    passwordSecret:
      secretName: mysecret
      password: mypassword
  config: 
    value.converter: io.confluent.connect.avro.AvroConverter 
    value.converter.schema.registry.url: https://mynamespace.servicebus.windows.net/myschemagroup
    value.converter.enhanced.avro.schema.support: true

    group.id: dev.mygroup.cdc.hub 
    offset.storage.topic: dev.mygroup.offsets.hub 
    config.storage.topic: dev.mygroup.configs.hub 
    status.storage.topic: dev.mygroup.status.hub 
    config.storage.replication.factor: -1
    offset.storage.replication.factor: -1
    status.storage.replication.factor: -1

  build:
    output:
      type: docker
      image: mycontainerregistry
      pushSecret: mysecret
    plugins:
      - name: mongodb
        artifacts:
          - type: jar
            url: https://repo1.maven.org/maven2/org/mongodb/kafka/mongo-kafka-connect/1.5.1/mongo-kafka-connect-1.5.1-all.jar 
          - type: jar
            url: https://repo1.maven.org/maven2/org/apache/avro/avro-1.10.2/avro-1.10.2.jar           
      - name: avro-converter
        artifacts:
          - type: zip
            url: https://d2p6pa21dvn84.cloudfront.net/api/plugins/confluentinc/kafka-connect-avro-converter/versions/7.7.2/confluentinc-kafka-connect-avro-converter-7.7.2.zip

我曾尝试将value.converter.schema.registry.url改为http://localhost:8081,依然失败,报错信息如下:

org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:230)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:156)
        at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.convertTransformedRecord(AbstractWorkerSourceTask.java:494)
        at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.sendRecords(AbstractWorkerSourceTask.java:402)
        at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:367)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259)
        at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:77)
        at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:237)
        at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
        at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
        at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
        at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
        at java.base/java.lang.Thread.run(Thread.java:840)
Caused by: org.apache.kafka.connect.errors.DataException: Failed to serialize Avro data from topic mytopic:
        at io.confluent.connect.avro.AvroConverter.fromConnectData(AvroConverter.java:107)
        at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.lambda$convertTransformedRecord$6(AbstractWorkerSourceTask.java:494)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:180)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:214)
        ... 13 more
Caused by: org.apache.kafka.common.errors.SerializationException: Error registering Avro schema"string"
        at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.toKafkaException(AbstractKafkaSchemaSerDe.java:917)
        at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:187)
        at io.confluent.connect.avro.AvroConverter$Serializer.serialize(AvroConverter.java:177)
        at io.confluent.connect.avro.AvroConverter.fromConnectData(AvroConverter.java:94)
        ... 16 more
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Unable to parse error message from schema registry: '(<Fault xmlns="http://schemas.microsoft.com/ws/2005/05/envelope/none"><Code><Value>Receiver</Value><Subcode><Value xmlns:a="http://schemas.microsoft.com/net/2005/12/windowscommunicationfoundation/dispatcher">a:InternalServiceFault</Value></Subcode></Code><Reason><Text xml:lang="en-US">The server was unable to process the request due to an internal error.  For more information about the error, either turn on IncludeExceptionDetailInFaults (either from ServiceBehaviorAttribute or from the <serviceDebug> configuration behavior) on the server in order to send the exception information back to the client, or turn on tracing as per the Microsoft .NET Framework SDK documentation and inspect the server trace logs.</Text></Reason></Fault>)'; error code: 50005
        at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:345)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:418)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:634)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:618)
        at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:610)
        at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.registerAndGetId(CachedSchemaRegistryClient.java:337)
        at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.registerWithResponse(CachedSchemaRegistryClient.java:458)
        at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.registerWithResponse(CachedSchemaRegistryClient.java:434)
        at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.registerWithResponse(AbstractKafkaSchemaSerDe.java:550)
        at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:120)
        ... 18 more
解决方案

问题根源在于Confluent Avro Converter与Azure Event Hub Schema Registry的API不兼容。Confluent的转换器默认针对Confluent Schema Registry设计,无法直接和Azure的Schema Registry交互,需要替换为Azure官方提供的转换器。

1. 替换Avro转换器并添加依赖

在KafkaConnect的build.plugins中移除Confluent的Avro Converter,添加Azure Schema Registry的Kafka Connect转换器依赖:

plugins:
  - name: mongodb
    artifacts:
      - type: jar
        url: https://repo1.maven.org/maven2/org/mongodb/kafka/mongo-kafka-connect/1.5.1/mongo-kafka-connect-1.5.1-all.jar 
      - type: jar
        url: https://repo1.maven.org/maven2/org/apache/avro/avro/1.10.2/avro-1.10.2.jar           
  - name: azure-schema-registry-converter
    artifacts:
      - type: jar
        url: https://repo1.maven.org/maven2/com/azure/messaging/eventhubs/azure-eventhubs-kafka-avro-converter/3.33.0/azure-eventhubs-kafka-avro-converter-3.33.0.jar
      - type: jar
        url: https://repo1.maven.org/maven2/com/azure/azure-core/1.49.0/azure-core-1.49.0.jar
      - type: jar
        url: https://repo1.maven.org/maven2/com/azure/azure-core-serializer-avro/1.15.0/azure-core-serializer-avro-1.15.0.jar
      - type: jar
        url: https://repo1.maven.org/maven2/com/azure/messaging/eventhubs/azure-messaging-eventhubs/5.17.0/azure-messaging-eventhubs-5.17.0.jar

2. 修改Kafka Connect配置

更新config中的转换器及Schema Registry相关参数,使用Azure的转换器类,并配置正确的认证信息:

config: 
  value.converter: com.azure.messaging.eventhubs.kafka.connect.avro.EventHubsAvroConverter
  value.converter.schema.registry.url: https://mynamespace.servicebus.windows.net/
  value.converter.schema.group.id: myschemagroup
  # 使用连接字符串认证(推荐)
  value.converter.connection.string: "${file:/opt/kafka/external-configuration/secrets/mysecret:connectionString}"
  # 或者使用AAD认证(如需)
  # value.converter.client.id: "<your-client-id>"
  # value.converter.tenant.id: "<your-tenant-id>"
  # value.converter.client.secret: "${file:/opt/kafka/external-configuration/secrets/mysecret:clientSecret}"

  group.id: dev.mygroup.cdc.hub 
  offset.storage.topic: dev.mygroup.offsets.hub 
  config.storage.topic: dev.mygroup.configs.hub 
  status.storage.topic: dev.mygroup.status.hub 
  config.storage.replication.factor: -1
  offset.storage.replication.factor: -1
  status.storage.replication.factor: -1

注意:

  • schema.registry.url只需指定Event Hub命名空间的基础URL,不需要包含schema group路径
  • schema.group.id单独配置为你的schema group名称
  • 认证优先使用Event Hub的连接字符串,需将连接字符串存入Kubernetes Secret中,通过文件引用方式加载

3. 调整认证配置

原配置中Event Hub的认证使用username: $MyConnectionString的方式可以保留,但需要确保Secret中的mypassword字段是正确的连接字符串内容。

4. 重新构建并部署Kafka Connect

应用上述配置后,重新构建镜像并部署Kafka Connect实例,Avro序列化即可正常与Azure Schema Registry交互。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 21:14:52