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"
}