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

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类型,两者无法直接映射。


修复步骤

  1. 补全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"
  1. 调整ValueToKey字段列表,仅保留分区键:
"transforms.ValueToKey.fields": "trackable_type,trackable_id"
  1. 修正Created_at类型匹配(二选一):
    • 方案一:修改Avro Schema,将created_at改为long类型(对应Cassandra TIMESTAMP的毫秒级时间戳),重新生成消息
    • 方案二:添加Cast转换,将string类型的created_at转为long:
"transforms": "ValueToKey,CastCreatedAt",
"transforms.CastCreatedAt.type": "org.apache.kafka.connect.transforms.Cast$Value",
"transforms.CastCreatedAt.spec": "created_at:long"
  1. 修正ID字段匹配(三选一):
    • 修改Cassandra表的id字段为TEXT类型
    • 修改Avro Schema的id为UUID类型(需确保生产者生成合法UUID)
    • 添加Cast转换,将string类型的id转为UUID:
"transforms": "ValueToKey,CastCreatedAt,CastId",
"transforms.CastId.type": "org.apache.kafka.connect.transforms.Cast$Value",
"transforms.CastId.spec": "id:uuid"

内容的提问来源于stack exchange,提问作者Rafael Simões

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:05:19