Kafka Connect ValueToKey转换异常:Cassandra复合主键配置故障
Cassandra Sink Connector复合主键配置问题排查
配置Cassandra Sink Connector时,尝试为Cassandra表设置复合主键,但遇到字段类型不匹配错误,生产者使用Avro序列化消息。
连接器配置
{ "connector.class": "io.confluent.connect.cassandra.CassandraSinkConnector", "tasks.max": "1", "topics": "events", "cassandra.contact.points": "cassandra", "cassandra.local.datacenter": "dc1", "cassandra.port": "9042", "cassandra.keyspace": "mykeyspace", "cassandra.table": "events_by_trackable", "cassandra.security": "NONE", "key.converter.schema.registry.url": "http://schema_registry:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://schema_registry:8081", "value.converter.schemas.enable": "true", "confluent.topic.replication.factor": 1, "confluent.topic.bootstrap.servers": "kafka_broker:9092", "transforms": "ValueToKey", "transforms.ValueToKey.type": "org.apache.kafka.connect.transforms.ValueToKey", "transforms.ValueToKey.fields": "trackable_type,trackable_id,created_at" }
Kafka Avro消费者输出
kafka-avro-console-consumer --bootstrap-server kafka_broker:9092 --property schema.registry.url=http://schema_registry:8081 --topic events --from-beginning --max-messages 10 --property print.key=true null {"id":"f3c183a2-3c0c-11ee-be54-da85ce8f1b40","actor_id":"123","trackable_owner_id":"12312","event_type":"View","trackable_type":"Product","trackable_id":"1234","created_at":1692173682188390,"metadata":{"user-agent":"Firefox"}}
错误信息
curl -s "http://localhost:8083/connectors/cassandra-sink/status" | jq -r ".tasks[0].trace" ... at java.base/java.lang.Thread.run(Thread.java:829) Caused by: org.apache.kafka.connect.errors.DataException: Exception thrown while processing one of the fields at io.confluent.connect.cassandra.BoundStatementConverter.convertStruct(BoundStatementConverter.java:406) at io.confluent.connect.cassandra.BoundStatementConverter.convert(BoundStatementConverter.java:293) at io.confluent.connect.cassandra.CassandraSinkTask.put(CassandraSinkTask.java:120) at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:590) ... 11 more Caused by: com.datastax.oss.driver.api.core.type.codec.CodecNotFoundException: Codec not found for requested operation: [BLOB <-> java.lang.String] at com.datastax.oss.driver.internal.core.type.codec.registry.CachingCodecRegistry.createCodec(CachingCodecRegistry.java:609) ...
Schema配置
curl http://localhost:8081/subjects/events/versions/latest {"subject":"events","version":1,"id":8,"schema":"{\"type\":\"record\",\"name\":\"events\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"actor_id\",\"type\":\"string\"},{\"name\":\"trackable_owner_id\",\"type\":\"string\"},{\"name\":\"event_type\",\"type\":\"string\"},{\"name\":\"trackable_type\",\"type\":\"string\"},{\"name\":\"trackable_id\",\"type\":\"string\"},{\"name\":\"created_at\",\"type\":\"string\"},{\"name\":\"metadata\",\"type\":{\"type\":\"map\",\"values\":\"string\"}}]}"}
Cassandra表结构
CREATE TABLE IF NOT EXISTS mykeyspace.events_by_trackable ( id UUID, actor_id TEXT, trackable_owner_id TEXT, event_type TEXT, trackable_type TEXT, trackable_id TEXT, created_at TIMESTAMP, metadata MAP<TEXT, TEXT>, PRIMARY KEY ((trackable_type, trackable_id), created_at) ) WITH CLUSTERING ORDER BY (created_at DESC);
问题分析与修复方案
核心问题1:Key Converter配置缺失
连接器仅指定了key.converter.schema.registry.url,但未设置key.converter。默认情况下Kafka Connect使用ByteArrayConverter处理消息键,导致ValueToKey转换后的结构化键被当成二进制Blob处理,与Cassandra的TEXT类型不匹配,触发[BLOB <-> java.lang.String]错误。
核心问题2:Created_at类型与用途错误
- Avro Schema中
created_at定义为string类型,但Cassandra表中是TIMESTAMP类型,类型不匹配 created_at是Cassandra的聚类键,不需要加入ValueToKey的字段列表(仅需将分区键trackable_type、trackable_id作为消息键)
核心问题3:ID字段类型不匹配
Avro中id是string类型,但Cassandra表中是UUID类型,两者无法直接映射。
修复步骤
- 补全Key Converter配置,与Value Converter保持一致:
"key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "http://schema_registry:8081", "key.converter.schemas.enable": "true"
- 调整ValueToKey字段列表,仅保留分区键:
"transforms.ValueToKey.fields": "trackable_type,trackable_id"
- 修正Created_at类型匹配(二选一):
- 方案一:修改Avro Schema,将
created_at改为long类型(对应Cassandra TIMESTAMP的毫秒级时间戳),重新生成消息 - 方案二:添加Cast转换,将string类型的
created_at转为long:
- 方案一:修改Avro Schema,将
"transforms": "ValueToKey,CastCreatedAt", "transforms.CastCreatedAt.type": "org.apache.kafka.connect.transforms.Cast$Value", "transforms.CastCreatedAt.spec": "created_at:long"
- 修正ID字段匹配(三选一):
- 修改Cassandra表的
id字段为TEXT类型 - 修改Avro Schema的
id为UUID类型(需确保生产者生成合法UUID) - 添加Cast转换,将string类型的
id转为UUID:
- 修改Cassandra表的
"transforms": "ValueToKey,CastCreatedAt,CastId", "transforms.CastId.type": "org.apache.kafka.connect.transforms.Cast$Value", "transforms.CastId.spec": "id:uuid"
内容的提问来源于stack exchange,提问作者Rafael Simões
相关产品推荐
相关产品推荐

