【问题标题】:Kafka connect possible to use custom query with bulk mode?Kafka连接可以使用批量模式的自定义查询吗?
【发布时间】:2020-09-27 10:29:28
【问题描述】:

我正在尝试发送 7 天前的每一行的记录。这是我正在处理的配置,但它 即使查询在数据库服务器上生成记录也不起作用。

{
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": 1,
    "mode": "bulk",
    "connection.url": "jdbc:mysql://mysql:3300/test_db?user=root&password=password",
    "query": "SELECT * FROM test_table WHERE DATEDIFF(CURDATE(), test_table.modified) = 7;",
    "topic.prefix": "test-jdbc-",
    "poll.interval.ms": 10000
}

【问题讨论】:

    标签: apache-kafka apache-kafka-connect


    【解决方案1】:

    JDBC 源连接器使用 JDBC 驱动程序将数据从关系数据库导入 Apache Kafka 主题。 数据会定期加载,要么根据时间戳递增,要么批量加载。最初,尽管在创建 JDBC 连接器时模式增量或批量加载,但它会在仅加载时间戳列上的新行或修改行之后将所有数据加载到主题中。

    批量:这种模式是未经过滤的,因此根本不是增量的。它将在每次迭代时从表中加载所有行。如果您想定期转储最终删除条目并且下游系统可以安全地处理重复项的整个表,这将很有用。 这意味着您不能使用批量模式增量加载过去 7 天

    时间戳列:在此模式下,包含修改时间戳的单个列用于跟踪上次处理数据的时间,并仅查询自该时间以来已修改的行。在这里您可以加载增量数据。但是当你第一次创建它时它是如何工作的,它将加载数据库表中所有可用的数据,因为对于 JDBC 连接器,这些是新数据。以后它只会加载新的或修改过的数据。

    现在根据您的要求,您似乎正在尝试以某个时间间隔加载所有数据,该时间间隔将配置为“poll.interval.ms”:10000。我看到您的 JDBC 连接设置符合定义,而查询可能不是工作尝试使用如下查询。似乎 JDBC 连接器将查询包装为一个表,如果添加 where 则该表不起作用。

    "query": "select * from (select * from test_table where  modified > now() - interval '7' day) o",
    

    试试下面的设置

    {
      "name": "application_name",
      "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
      "key.converter": "org.apache.kafka.connect.json.JsonConverter",
      "value.converter": "org.apache.kafka.connect.json.JsonConverter",
      "connection.url": "jdbc:mysql://mysql:3300/test_db",
      "connection.user": "root",
      "connection.password": "password",
      "connection.attempts": "1",
      "mode": "bulk",
      "validate.non.null": false,
      "query": "select * from (select * from test_table where  modified > now() - interval '7' day) o",
      "table.types": "TABLE",
      "topic.prefix": "test-jdbc-",
     "poll.interval.ms": 10000
      "schema.ignore": true,
      "key.converter.schemas.enable": "false",
      "value.converter.schemas.enable": "false"
      
    }
    

    【讨论】:

    • 非常感谢!我被困了一段时间。
    猜你喜欢
    • 2019-08-08
    • 2021-07-10
    • 1970-01-01
    • 2020-04-09
    • 2020-07-27
    • 1970-01-01
    • 2021-04-07
    • 2021-10-24
    • 1970-01-01
    相关资源
    最近更新 更多