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 &lt;serviceDebug&gt; 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

