Kafka Connect同步Kafka数据至Cassandra表时添加插入/更新时间戳的实现方案及JSON插入now()函数报错咨询
首先咱们先理清你遇到的报错根源:Cassandra的JSON插入API要求传入的是标准JSON格式,而now()是CQL的内置函数,不属于JSON的合法字面量(JSON只认字符串、数字、布尔值、null、对象、数组),所以解析器会报错无法识别now()这个token,这就是你看到那个错误的核心原因。
接下来给你几个无需修改Cassandra集群代码的可行方案:
方案1:用Kafka Connect的SMT(Single Message Transform)添加时间戳
这是最推荐的方案,直接在Kafka Connect的Sink Connector配置里加入消息转换逻辑,在数据到达Cassandra前自动插入当前时间戳字段。
你可以用官方提供的InsertField SMT,配置示例如下:
# 启用SMT并命名为addTimestamp transforms=addTimestamp # 指定SMT为Value级别的插入字段操作 transforms.addTimestamp.type=org.apache.kafka.connect.transforms.InsertField$Value # 设置要添加的时间戳字段名,和Cassandra表字段保持一致 transforms.addTimestamp.timestamp.field=last_updated_on_cassandra # 定义时间戳格式,Cassandra的TIMESTAMP类型兼容该格式 transforms.addTimestamp.timestamp.format=yyyy-MM-dd'T'HH:mm:ss.SSSZ
配置完成后,Kafka Connect会自动在每条消息的Value中添加last_updated_on_cassandra字段,值为Connector处理这条消息时的当前时间,此时再用JSON插入,该字段就是合法的字符串时间戳,Cassandra能正常解析。
方案2:给Cassandra字段设置默认值
如果你的业务只需要记录插入时间(更新时不需要刷新时间戳),可以在创建Cassandra表时直接给目标字段设置默认值:
CREATE TABLE keyspace1.table1 ( lhlt TEXT PRIMARY KEY, last_updated_on_cassandra TIMESTAMP DEFAULT toTimestamp(now()) );
这样当Kafka Connect发送的JSON不包含该字段时,Cassandra会自动用当前时间填充。但要注意:如果是基于主键的Upsert更新操作,若消息里未指定这个字段,Cassandra不会更新它的值,只会保留最初的插入时间。
方案3:自定义CQL插入语句(绕过JSON模式)
如果你想直接在插入时使用CQL的now()函数,可以修改Cassandra Sink Connector的配置,指定自定义的插入语句,放弃默认的JSON插入模式。
以DataStax的Cassandra Sink Connector为例,配置示例如下:
# 指定自定义CQL查询语句 cassandra.query=INSERT INTO keyspace1.table1 (lhlt, last_updated_on_cassandra) VALUES (?, toTimestamp(now())) # 配置字段映射,将Kafka消息中的lhlt字段对应到查询的第一个参数 cassandra.mapping=lhlt=0
这种方式下,Connector会直接执行你指定的CQL语句,把消息中的lhlt字段作为参数传入,第二个参数用toTimestamp(now())生成当前时间戳,完全避开了JSON解析的问题。
针对你报错的直接解决方法
如果坚持要用JSON插入模式,必须把now()替换成实际的时间戳字符串(比如"2024-05-20T12:34:56.789+0000"),而不是CQL函数。上面的方案1就是自动完成这个转换的最优方式。
内容的提问来源于stack exchange,提问作者Otto

