【问题标题】:Confluent: ERROR Failed to run query for table TimestampIncrementingTableQuerier mysql-jdbcConfluent:错误无法对表TimestampIncrementingTableQuerier mysql-jdbc运行查询
【发布时间】:2019-05-01 12:04:58
【问题描述】:

我正在尝试在 MySQL 中使用模式时间戳,由于我的表大小为 2.6 GB,因此行数有限。

这是我正在使用的连接器属性:

{
        "name": "jdbc_source_mysql_registration_query",
        "config": {
                 "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
                 "key.converter": "io.confluent.connect.avro.AvroConverter",
                 "key.converter.schema.registry.url": "http://localhost:8081",
                 "value.converter": "io.confluent.connect.avro.AvroConverter",
                 "value.converter.schema.registry.url": "http://localhost:8081",
                 "connection.url": "jdbc:mysql://localhost:3310/users?zeroDateTimeBehavior=ROUND&useCursorFetch=true&defaultFetchSize=1000&user=kotesh&password=kotesh",
                 "query": "SELECT matriid,DateUpdated  from users.employee WHERE date(DateUpdated)>='2018-11-28' ",
                 "mode": "timestamp",
                 "timestamp.column.name": "DateUpdated",
                 "validate.non.null": "false",
                 "topic.prefix": "mysql-prod-kot-"
        }
}

我得到如下:

INFO TimestampIncrementingTableQuerier{table=null, query='SELECT matriid,DateUpdated from users.employee WHERE date(DateUpdated)>='2018-11-28'', topicPrefix='mysql-prod-kot-', incrementingColumn='', timestampColumns=[DateUpdated]} 准备好的 SQL 查询:SELECT matriid,DateUpdated from users.employee WHERE 日期(更新日期)>='2018-11-28'DateUpdated > ?和 DateUpdated DateUpdated ASC (io.confluent.connect.jdbc.source.TimestampIncrementingTableQuerier:161) [2018-11-29 17:29:00,981] 错误无法对表 TimestampIncrementingTableQuerier{table=null, query='SELECT 运行查询 matriid,DateUpdated from users.employee WHERE date(DateUpdated)>='2018-11-28'', topicPrefix='mysql-prod-kot-', incrementingColumn='',timestampColumns=[DateUpdated]}:{} (io.confluent.connect.jdbc.source.JdbcSourceTask:328) java.sql.SQLSyntaxErrorException:您的 SQL 语法有错误;检查与您的 MySQL 服务器版本相对应的手册 在 'WHERE DateUpdated > '1970-01-01 附近使用正确的语法 00:00:00.0' AND DateUpdated

【问题讨论】:

    标签: jdbc apache-kafka apache-kafka-connect confluent-platform confluent-schema-registry


    【解决方案1】:

    报错如图:

    java.sql.SQLSyntaxErrorException: You have an error in your SQL syntax; 
    check the manual that corresponds to your MySQL server version for the right syntax to use near 
    'WHERE `DateUpdated` > '1970-01-01 00:00:00.0' AND `DateUpdated` < '2018-11-29 17' at line 1
    

    这是因为您使用的是query,但也使用了"mode": "timestamp",因此当您还在查询中指定了一个WHERE 子句时,连接器会尝试附加它自己的子句,这会导致SQL 无效

    对于 JDBC 源连接器,每个 docs

    为了正确构造增量查询,必须可以在该查询中附加 WHERE 子句(即不得使用 WHERE 子句)。如果使用 WHERE 子句,它必须自己处理增量查询。

    【讨论】:

    • 跟进问题 - 您如何指定连接器的起始时间戳
    【解决方案2】:

    这是因为您尝试同时使用"mode": "timestamp"queryTimestampIncrementingTableQuerierWHERE 子句附加到与query 中现有WHERE 子句冲突的查询。

    JDBC source connector docs 对此很清楚:

    query

    如果指定,则执行查询以选择新的或更新的行。采用 如果要连接表,请选择此设置中的列子集 表格或过滤数据。如果使用,此连接器将仅复制数据 使用这个查询——全表复制将被禁用。不同的 查询模式仍可用于增量更新,但为了 正确构造增量查询,必须可以 将 WHERE 子句附加到此查询(即,可能没有 WHERE 子句 用过的)。 如果使用 WHERE 子句,它必须处理增量查询 本身

    作为一种解决方法,您可以将查询修改为(取决于您使用的 SQL 风格)

    SELECT * FROM ( SELECT * FROM table WHERE ...)
    

    WITH a AS
       SELECT * FROM b
        WHERE ...
    SELECT * FROM a
    

    例如,在您的情况下,查询应该是

    "query":"SELECT * FROM (SELECT matriid,DateUpdated  from users.employee WHERE date(DateUpdated)>='2018-11-28') o"
    

    【讨论】:

      猜你喜欢
      • 2012-05-27
      • 2014-09-23
      • 2020-09-04
      • 1970-01-01
      • 2018-10-25
      • 1970-01-01
      • 2018-12-19
      • 1970-01-01
      • 2019-03-03
      相关资源
      最近更新 更多