使用Landoop CassandraSinkConnector同步Kafka到Cassandra遇字段错误
解决CassandraSinkConnector同步Kafka数据到Cassandra的FIELD_ID字段无效错误
结合你给出的报错日志、Avro Schema和Sink配置,咱们来一步步定位并解决这个问题。
首先看你碰到的核心报错:
Caused by: java.lang.IllegalArgumentException: A KCQL error occurred.FIELD_ID is not a valid field name at com.datamountaineer.streamreactor.connect.converters.Transform$.raiseException$1(Transform.scala:40) at com.datamountaineer.streamreactor.connect.converters.Transform$.apply(Transform.scala:83) at com.datamountaineer.streamreactor.connect.cassandra.sink.CassandraJsonWriter$$anonfun$com$datamountaineer$streamreactor$connect$cassandra$sink$CassandraJsonWriter$$insert$1.apply(CassandraJsonWriter.scala:182) at com.datamountaineer.streamreactor.connect.cassandra.sink.CassandraJsonWriter$$anonfun$com$datamountaineer$streamreactor$connect$cassandra$sink$CassandraJsonWriter$$insert$1.apply(CassandraJsonWriter.scala:181) at scala.collection.immutable.Map$Map1.foreach(Map.scala:116)
再看你的Avro Schema,数据是嵌套结构:外层的data字段是一个Data类型的record,里面才包含FIELD_ID和SERVER_INSTANCE这两个业务字段:
{ "type": "record", "name": "DataRecord", "namespace": "com.attunity.queue.msg.test 1 dlx express.DCSDBA.PURGE_SETUP", "fields": [ { "name": "data", "type": { "type": "record", "name": "Data", "fields": [ { "name": "FIELD_ID", "type": ["null", "string"], "default": null }, { "name": "SERVER_INSTANCE", "type": ["null", "int"], "default": null } ] } }, { "name": "beforeData", "type": ["null", "Data"], "default": null }, { "name": "headers", "type": { "type": "record", "name": "Headers", "namespace": "com.attunity.queue.msg", "fields": [ { "name": "operation", "type": { "type": "enum", "name": "operation", "symbols": ["INSERT", "UPDATE", "DELETE", "REFRESH"] } }, { "name": "changeSequence", "type": "string" }, { "name": "timestamp", "type": "string" }, { "name": "streamPosition", "type": "string" }, { "name": "transactionId", "type": "string" }, { "name": "changeMask", "type": ["null", "bytes"], "default": null }, { "name": "columnMask", "type": ["null", "bytes"], "default": null }, { "name": "transactionEventCounter", "type": ["null", "long"], "default": null }, { "name": "transactionLastEvent", "type": ["null", "boolean"], "default": null } ] } } ] }
你的Sink配置里的KCQL语句是:
connect.cassandra.kcql=INSERT INTO purge_setup SELECT data.* FROM DCSDBA.PURGE_SETUP1
问题根源分析
Cassandra字段名大小写冲突
Cassandra默认会把所有未用双引号包裹的字段名自动转为小写。如果你的purge_setup表是用小写字段名(比如field_id)创建的,那么Connector通过data.*提取出的大写FIELD_ID在Cassandra中找不到,就会抛出“不是有效字段名”的错误。嵌套字段的KCQL解析问题
Landoop的CassandraSinkConnector对嵌套Avro结构的*.语法支持有限,无法自动处理嵌套字段的大小写映射,也可能无法正确识别嵌套层级的字段。
解决方案
修改KCQL配置,明确指定字段的映射关系,同时适配Cassandra的小写命名规范:
connect.cassandra.kcql=INSERT INTO purge_setup SELECT data.FIELD_ID AS field_id, data.SERVER_INSTANCE AS server_instance FROM DCSDBA.PURGE_SETUP1
如果你的Cassandra表确实是用大写字段名创建的(创建时需要用双引号包裹,比如CREATE TABLE purge_setup ("FIELD_ID" text, ...)),那可以保持字段名一致,但这种写法不符合Cassandra的最佳实践,不推荐。
另外还要注意两个细节:
- 检查Cassandra表的字段类型是否和Kafka消息中的类型匹配:比如
FIELD_ID对应字符串类型,SERVER_INSTANCE对应整数类型。 - 你的Sink配置中
topics写的是DCSDBA.PURGE_SETUP1,但Schema对应的Topic是DCSDBA.PURGE_SETUP,如果这是笔误,也需要修正,否则Connector会读取不到正确的Topic数据。
内容的提问来源于stack exchange,提问作者Otto
相关产品推荐
相关产品推荐

