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

DataStax Cassandra Sink Connector按条件从Kafka主题摄入数据问题

问题修复与方案实现

一、当前单表过滤配置错误修复

你配置报错的核心原因是语法使用错误:DataStax Cassandra Sink 连接器的自定义查询中,仅INSERT取值位的字段需要加冒号前缀,WHERE条件中直接引用字段名即可,不需要额外加冒号,也不需要写FROM 主题名的SELECT语法,正确配置如下:

{
  "name": "cassandra-json-sink",
  "config": {
    "connector.class": "com.datastax.oss.kafka.sink.CassandraSinkConnector",
    "tasks.max": "1",
    "topics": "json_test_topic",
    "contactPoints": "cassandra",
    "loadBalancing.localDc": "datacenter1",
    "port": 9042,
    "auth.username": "cassandra",
    "auth.password": "cassandra",
    "topic.json_test_topic.kconnect_json.customer.mapping": "id=key, name=value.name, lname=value.lname, adress=value.adress",
    "topic.json_test_topic.kconnect_json.customer.query": "INSERT INTO kconnect_json.customer(id, name, lname, adress) VALUES (:id, :name, :lname, :adress) WHERE name = 'john';",
    "topic.json_test_topic.kconnect_json.customer.deletesEnabled": false,
    "dropInvalidMessage": true,
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "key.converter.schemas.enable": false,
    "value.converter.schemas.enable": false
  }
}

新增的dropInvalidMessage配置用于将不匹配条件的消息直接丢弃,避免触发异常导致任务退出。

二、单主题路由到多Cassandra表实现

你需要的按消息头分表的需求完全可以在同一个连接器配置中实现,不需要部署多个连接器实例。配置逻辑为给每个目标表单独配置映射规则和匹配查询条件即可,示例配置如下(假设消息头中判断字段为event_type,对应值为A、B、C,目标表为table_a、table_b、table_c):

{
  "name": "cassandra-json-sink-multi-table",
  "config": {
    "connector.class": "com.datastax.oss.kafka.sink.CassandraSinkConnector",
    "tasks.max": "1",
    "topics": "json_test_topic",
    "contactPoints": "cassandra",
    "loadBalancing.localDc": "datacenter1",
    "port": 9042,
    "auth.username": "cassandra",
    "auth.password": "cassandra",
    // 表A配置
    "topic.json_test_topic.kconnect_json.table_a.mapping": "id=key, col1=value.col1, col2=value.col2",
    "topic.json_test_topic.kconnect_json.table_a.query": "INSERT INTO kconnect_json.table_a(id, col1, col2) VALUES (:id, :col1, :col2) WHERE header.event_type = 'A';",
    "topic.json_test_topic.kconnect_json.table_a.deletesEnabled": false,
    // 表B配置
    "topic.json_test_topic.kconnect_json.table_b.mapping": "id=key, field1=value.field1, field2=value.field2",
    "topic.json_test_topic.kconnect_json.table_b.query": "INSERT INTO kconnect_json.table_b(id, field1, field2) VALUES (:id, :field1, :field2) WHERE header.event_type = 'B';",
    "topic.json_test_topic.kconnect_json.table_b.deletesEnabled": false,
    // 表C配置
    "topic.json_test_topic.kconnect_json.table_c.mapping": "id=key, item1=value.item1, item2=value.item2",
    "topic.json_test_topic.kconnect_json.table_c.query": "INSERT INTO kconnect_json.table_c(id, item1, item2) VALUES (:id, :item1, :item2) WHERE header.event_type = 'C';",
    "topic.json_test_topic.kconnect_json.table_c.deletesEnabled": false,
    "dropInvalidMessage": true,
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "key.converter.schemas.enable": false,
    "value.converter.schemas.enable": false
  }
}

注意事项

  • WHERE条件中可以直接引用三类字段:消息key字段(直接写字段名)、消息value字段(直接写字段名)、消息头字段(格式为header.头字段名)
  • 所有过滤逻辑由连接器侧完成,不会向Cassandra发送非法CQL语句

内容的提问来源于stack exchange,提问作者H.Demir

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 19:39:03