使用Confluent JDBC连接器同步Snowflake到Kafka时遇标识符错误
问题排查与解决方案
核心问题分析
从错误日志、表结构和连接器配置来看,存在两个关键问题导致SQL编译错误:
1. 查询表名不匹配
日志中显示连接器执行的查询为 select * from oslo_cycle_trip,但你实际创建的表是 cycle_trip。尽管连接器配置中指定了 query=select * from cycle_trip,但日志中的表名明显不符,说明配置可能未正确生效(比如未重启连接器、存在旧配置残留)。
2. Snowflake标识符大小写冲突
Snowflake的标识符规则:创建表/列时如果未使用双引号,会自动将标识符转为大写。你创建表时的列定义是 id integer(无引号),因此实际列名是 ID,而非小写的 id。但连接器配置中 incrementing.column.name=id,JDBC驱动会将其转义为带双引号的 "id",而Snowflake中不存在这个小写列名,最终导致报错invalid identifier ""id""。
具体修复步骤
步骤1:修正查询表名并确保配置生效
- 确认连接器配置中的
query参数确实为select * from cycle_trip,无拼写错误 - 重启Kafka Connect集群或该连接器实例,确保新配置覆盖旧配置
步骤2:解决标识符大小写问题(二选一即可)
方案A:修改连接器配置的增量列名为大写
将配置中的 incrementing.column.name=id 修改为:
incrementing.column.name=ID
方案B:在查询语句中显式指定列并别名(如需保持小写列名)
修改 query 参数为:
query=select "ID" as id, started_at, ended_at, duration from cycle_trip
这样既匹配Snowflake中的实际列名,又能让连接器识别小写的id作为增量列。
额外验证建议
- 直接在Snowflake客户端执行日志中的报错查询语句(
select * from oslo_cycle_trip),快速确认表名是否存在问题 - 执行
DESCRIBE TABLE cycle_trip;查看列名的实际大小写,验证是否为ID
内容的提问来源于stack exchange,提问作者user565
相关产品推荐
相关产品推荐

