Couchbase-Kafka-SQL Server数据同步Avro Schema适配问题求助
需要实现流程:Couchbase创建文档 → Kafka接收消息 → SQL Server将消息存储到表中。当前Kafka连接器配置如下,创建Couchbase文档时Sink连接器立即失败,排查发现是Couchbase不支持Avro导致配置不兼容,寻求纯Kafka UI配置的解决方案。
SINK连接器配置
{ "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "table.name.format": "mytopic", "connection.password": "******", "tasks.max": "1", "topics": "usertopic", "schema.registry.url": "http://schema-registry:8081", "value.converter.schema.registry.url": "http://schema-registry:8081", "auto.evolve": "true", "connection.user": "Allen", "value.converter.schemas.enable": "true", "name": "usertopic", "auto.create": "true", "value.converter": "io.confluent.connect.avro.AvroConverter", "connection.url": "jdbc:sqlserver://192.168.0.1:1433;databaseName=pubs", "insert.mode": "insert", "pk.mode": "none" }
SOURCE连接器配置
{ "connector.class": "com.couchbase.connect.kafka.CouchbaseSourceConnector", "couchbase.persistence.polling.interval": "100ms", "tasks.max": "2", "couchbase.seed.nodes": "192.168.0.1", "couchbase.source.handler": "com.couchbase.connect.kafka.handler.source.RawJsonSourceHandler", "value.converter.schema.registry.url": "http://schema-registry:8081", "couchbase.bucket": "Embedding", "couchbase.username": "Admin", "name": "usertopic_source", "value.converter.schemas.enable": "true", "couchbase.password": "******", "couchbase.event.filter": "com.couchbase.connect.kafka.filter.AllPassFilter", "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "couchbase.topic": "usertopic" }
错误日志
org.apache.kafka.connect.errors.ConnectException: 错误处理器中超出容错阈值
at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:260)
at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:179)
at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:540)
at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:517)
at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:343)
at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:246)
at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:215)
at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:225)
at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:280)
at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:237)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source)
at java.base/java.util.concurrent.FutureTask.run(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
at java.base/java.lang.Thread.run(Unknown Source)
Caused by: org.apache.kafka.connect.errors.DataException: 无法将主题usertopic的数据反序列化为Avro:
at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:148)
at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$4(WorkerSinkTask.java:540)
at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:207)
at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:244)
... 14个更多错误
Caused by: org.apache.kafka.common.errors.SerializationException: 未知的魔术字节!
at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.getByteBuffer(AbstractKafkaSchemaSerDe.java:641)
at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer$DeserializationContext.(AbstractKafkaAvroDeserializer.java:382)
at io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer.deserializeWithSchemaAndVersion(AbstractKafkaAvroDeserializer.java:257)
at io.confluent.connect.avro.AvroConverter$Deserializer.deserialize(AvroConverter.java:199)
at io.confluent.connect.avro.AvroConverter.toConnectData(AvroConverter.java:126)
... 17个更多错误
核心问题是Source连接器输出的是原始JSON字节,而Sink连接器试图用Avro格式解析,导致序列化不兼容。只需调整两个连接器的转换器配置,统一使用JSON格式即可,无需依赖Schema Registry或自定义代码。
1. 修改SOURCE连接器配置
将值转换器改为JSON转换器,移除不必要的Schema Registry配置:
{ "connector.class": "com.couchbase.connect.kafka.CouchbaseSourceConnector", "couchbase.persistence.polling.interval": "100ms", "tasks.max": "2", "couchbase.seed.nodes": "192.168.0.1", "couchbase.source.handler": "com.couchbase.connect.kafka.handler.source.RawJsonSourceHandler", "couchbase.bucket": "Embedding", "couchbase.username": "Admin", "name": "usertopic_source", "value.converter.schemas.enable": "true", "couchbase.password": "******", "couchbase.event.filter": "com.couchbase.connect.kafka.filter.AllPassFilter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "couchbase.topic": "usertopic" }
关键修改点:
- 把
value.converter从ByteArrayConverter替换为org.apache.kafka.connect.json.JsonConverter,直接将Couchbase的JSON文档序列化为Kafka支持的JSON格式 - 移除
value.converter.schema.registry.url,因为JSON转换器不需要Schema Registry
2. 修改SINK连接器配置
同样将值转换器改为JSON转换器,移除Avro相关配置:
{ "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "table.name.format": "mytopic", "connection.password": "******", "tasks.max": "1", "topics": "usertopic", "auto.evolve": "true", "connection.user": "Allen", "value.converter.schemas.enable": "true", "name": "usertopic", "auto.create": "true", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "connection.url": "jdbc:sqlserver://192.168.0.1:1433;databaseName=pubs", "insert.mode": "insert", "pk.mode": "none", "json.schema.infer": "true" }
关键修改点:
- 把
value.converter从AvroConverter替换为org.apache.kafka.connect.json.JsonConverter,匹配Source的输出格式 - 移除
schema.registry.url和value.converter.schema.registry.url - 添加
json.schema.infer": "true",让连接器自动从JSON数据推断表结构(如果value.converter.schemas.enable设为false,则必须开启这个配置)
验证步骤
- 先删除原有的Source和Sink连接器
- 分别创建修改后的Source和Sink连接器
- 在Couchbase中创建测试文档,检查SQL Server的
mytopic表是否自动创建并插入数据
如果遇到表结构不匹配的问题,可以调整auto.evolve和json.schema.infer的配置,或者提前在SQL Server中创建对应结构的表,关闭auto.create。
内容的提问来源于stack exchange,提问作者user284331

