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
相关产品推荐
相关产品推荐

