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

使用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

问题根源分析

  1. Cassandra字段名大小写冲突
    Cassandra默认会把所有未用双引号包裹的字段名自动转为小写。如果你的purge_setup表是用小写字段名(比如field_id)创建的,那么Connector通过data.*提取出的大写FIELD_ID在Cassandra中找不到,就会抛出“不是有效字段名”的错误。

  2. 嵌套字段的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:58:38