【发布时间】:2019-11-09 04:45:27
【问题描述】:
我在 kafka 中有一个 JDBCSourceConnector,它使用查询从数据库中流式传输数据。 但我为选择数据而编写的查询有问题。
我在 Postgres psql 和 DBeaver 中测试了查询。它工作正常,但在 kafka 配置中,它会产生 SQL 语法错误
错误
错误无法运行表 TimestampIncrementingTableQuerier{name='null', query='select "Users".* from "Users" join "SchoolUserPivots" on "Users".id = "SchoolUserPivots".user_id where school_id = 1 and role_id = 2', topicPrefix='teacher', timestampColumn='"Users".updatedAt', incrementingColumn='id'}: {} (io.confluent.connect.jdbc.source.JdbcSourceTask:221) org.postgresql.util.PSQLException:错误:“WHERE”或附近的语法错误
配置 json
{
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"timestamp.column.name": "\"Users\".updatedAt",
"incrementing.column.name": "id",
"connection.password": "123",
"tasks.max": "1",
"query": "select \"Users\".* from \"Users\" join \"SchoolUserPivots\" on \"Users\".id = \"SchoolUserPivots\".user_id where school_id = 1 and role_id = 2",
"timestamp.delay.interval.ms": "5000",
"mode": "timestamp+incrementing",
"topic.prefix": "teacher",
"connection.user": "user",
"name": "SourceTeacher",
"connection.url": "jdbc:postgresql://ip:5432/school",
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter": "org.apache.kafka.connect.json.JsonConverter"
}
【问题讨论】:
-
删除列名周围的单引号:
where school_id = 1 and role_id = 2(而不是'shool_id') -
@a_horse_with_no_name 一样的错误,没有区别
-
用您的固定代码和新错误更新您的问题。我的猜测是,也许你修复了
school_id,但没有将该修复应用于role_id -
@MarkRotteveel 代码和错误已更新
-
您可以尝试将其更改为
select * from (<your original query) a。我猜卡夫卡正在添加自己的 where 子句。
标签: postgresql jdbc apache-kafka