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

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,则必须开启这个配置)

验证步骤

  1. 先删除原有的Source和Sink连接器
  2. 分别创建修改后的Source和Sink连接器
  3. 在Couchbase中创建测试文档,检查SQL Server的mytopic表是否自动创建并插入数据

如果遇到表结构不匹配的问题,可以调整auto.evolve和json.schema.infer的配置,或者提前在SQL Server中创建对应结构的表,关闭auto.create。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 12:16:00