【问题标题】:Does the kafka JDBC MySQL source connector need to have MySQL Server on localhost?kafka JDBC MySQL 源连接器是否需要在 localhost 上有 MySQL Server?
【发布时间】:2019-06-22 05:42:17
【问题描述】:

我对 Kafka 很陌生,我正在尝试使用 MySQL 源连接器和 Elasticsearch + Elastic 搜索接收器连接器启动并运行一个简单的 kafka 连接系统;用于基本数据流目的。

我正在按照以下步骤操作 https://www.confluent.io/blog/simplest-useful-kafka-connect-data-pipeline-world-thereabouts-part-1/ 及其第 2 部分 (我已经通过在源端有一个简单的生产者来验证 ES 的工作原理。)

除 MySQL 源连接器外,一切都按预期进行配置和工作。 我正在尝试所有这些的 VM 上没有安装 MySQL 服务器。本教程的 DBMS 部分我正在使用客户端来创建/更改和使用表。 因此,在源属性中,我尝试了:

"connection.url": "jdbc:mysql://IPaddressofDB:3306/DBname?user=uname&password=pwd"
"table.whitelist": "tablename"

要启动连接器,我只是做了一个./confluent load connector-name

一旦我加载源连接器并检查它的状态,它就会给出一个错误

"org.apache.kafka.connect.errors.ConnectException: Failed trying to validate that columns used for offsets are NOT NULL\n\t ...
 Caused by: java.sql.SQLSyntaxErrorException: Table 'admin_portal.tablename' doesn't exist\n\t
  1. 这是否正确?我完全错过了什么吗?

  2. 如何为我正在尝试的情况指定 connection.url:您尝试在哪里连接到不同的数据库服务器?几乎所有示例/git 问题等似乎都只指定 localhost。

  3. 我不确定admin_portal 来自哪里,我根本没有指定任何地方

****为@robin-moffat 的建议编辑(似乎给出与以前相同的错误)

sourceconfig.json:

{
        "name": "jdbc_source_mysql_new",
        "config": {
                "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
                "connection.url": "jdbc:mysql://ipaddress:3306/dbname?user=uname&password=pwd",
                "table.whitelist": "dbname.tablename",
                "topic.prefix": "mysql-new-",
                "mode":"incrementing",
                "incrementing.column.name": "colname"
                }
}

加载连接器:

>curl -X POST -H "Content-Type: application/json" --data @sourceconfig.json http://localhost:8083/connectors

检查连接器的状态:

>curl -X GET localhost:8083/connectors/jdbc_source_mysql_new/tasks/0/status

  {"state":"FAILED","
     "trace": 
     "org.apache.kafka.connect.errors.ConnectException: Failed trying to validate that columns used for offsets are NOT NULL\n\t
     at io.confluent.connect.jdbc.source.JdbcSourceTask.validateNonNullable(JdbcSourceTask.java:400)\n\t
     at io.confluent.connect.jdbc.source.JdbcSourceTask.start(JdbcSourceTask.java:156)\n\t
     at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:198)\n\t
     at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)\n\t
     at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)\n\t
     at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)\n\t
     at java.util.concurrent.FutureTask.run(FutureTask.java:266)\n\t
     at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)\n\t
     at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)\n\t
 at java.lang.Thread.run(Thread.java:748)\n

 Caused by: java.sql.SQLSyntaxErrorException: Table 'admin_portal.tablename' doesn't exist\n\t
 at com.mysql.cj.jdbc.exceptions.SQLError.createSQLException(SQLError.java:120)\n\t
 at com.mysql.cj.jdbc.exceptions.SQLError.createSQLException(SQLError.java:97)\n\t
 at com.mysql.cj.jdbc.exceptions.SQLExceptionsMapping.translateException(SQLExceptionsMapping.java:122)\n\t
 at com.mysql.cj.jdbc.StatementImpl.executeQuery(StatementImpl.java:1218)\n\t
 at com.mysql.cj.jdbc.DatabaseMetaData$7.forEach(DatabaseMetaData.java:2950)\n\t
 at com.mysql.cj.jdbc.DatabaseMetaData$7.forEach(DatabaseMetaData.java:2938)\n\t
 at com.mysql.cj.jdbc.IterateBlock.doForAll(IterateBlock.java:56)\n\t
 at com.mysql.cj.jdbc.DatabaseMetaData.getPrimaryKeys(DatabaseMetaData.java:2991)\n\t
 at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.primaryKeyColumns(GenericDatabaseDialect.java:696)\n\t
 at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.describeColumns(GenericDatabaseDialect.java:533)\n\t
 at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.describeColumns(GenericDatabaseDialect.java:513)\n\t
 at io.confluent.connect.jdbc.source.JdbcSourceTask.validateNonNullable(JdbcSourceTask.java:369)\n\t... 9 more\n",}

【问题讨论】:

  • 这与本地主机无关:错误表明您已连接(否则它永远不会产生SQLSyntaxErrorException)。布局并不是 kafka 所期望的。
  • 您好,您能否说明一下此布局中的意外情况?
  • 该错误表明 kafka 正在等待您的表中不存在的表。为什么会这样,不知道。

标签: mysql jdbc apache-kafka apache-kafka-connect


【解决方案1】:

kafka JDBC MySQL 源连接器是否需要在 localhost 上有 MySQL Server?

没有。它使用可以连接到远程实例上的服务器的 JDBC。

  1. 这是否正确?我完全错过了什么吗?

根据您的描述,您是正确的:)

  1. 如何为我正在尝试的情况指定 connection.url:您尝试连接到不同数据库服务器的位置?几乎所有示例/git 问题等似乎都只指定 localhost。

你可以看到an example here

您需要正确配置JDBC URL,can be found here for MySQL的语法。

  1. 我不确定 admin_portal 来自哪里,我根本没有指定任何地方

这将取决于您连接到数据库的用户的权限。您需要确保它可以访问要从中读取数据的表。您还可以限定您的表名,例如

"table.whitelist": "schema.tablename"

【讨论】:

  • 谢谢@robin。让我马上检查一下。所以,如果我理解正确的话,我只需要将 kafka-connect 的 REST API 用于远程服务器,在 localhost 的情况下就不需要了。
  • 嗨@Robin,感谢您的建议。但是我仍然遇到同样的错误。我已经编辑了问题以显示我尝试过的内容。我也先尝试了批量模式,但遇到了同样的错误
【解决方案2】:

在我将 My SQL 连接器版本从 8.x 降级到 5.1.47 并将其放置在正确的 $CLASSPATH 中后它工作了

mysql-connector-java-5.1.47.jar

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-08-30
    • 2017-10-12
    • 2022-06-10
    • 1970-01-01
    • 2020-01-15
    • 2012-11-12
    • 2018-09-01
    • 2011-06-28
    相关资源
    最近更新 更多