【问题标题】:Kafka connect JDBC sink isse with table name containing a " . "Kafka 使用包含“.”的表名连接 JDBC 接收器 isse
【发布时间】:2020-05-28 07:19:25
【问题描述】:

我正在尝试构建一个 kafka 连接 jdbc 接收器连接器。问题是,数据库表名包含一个点,并且在创建连接器时,该过程将表名拆分为两个导致未找到的数据库表。我尝试了多种方法来转义点,以便可以将其作为表名中的字符串读取,但没有任何效果..

这是真正的名字:

"table.name.format":"Bte3_myname.centrallogging",

这是错误:

原因:org.apache.kafka.connect.errors.ConnectException:表\"Bte3_myname\".\"centrallogrecord\"丢失。

这是我的配置文件:

{
    "name": "jdbc-connect-central-logging-sink",
    "config": 
    {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "3",
        "topics": "central_logging",

        "connection.url": "...",         
        "connection.user": "...",
        "connection.password": "...",
        "table.name.format":"Bte3_myname.centrallogging",
        "pk.mode": "kafka",

        "auto.create": "false",
        "auto.evolve": "false"
    }
}

有人知道如何在配置文件中正确解析吗?

非常感谢!

【问题讨论】:

  • Bte3_myname 可能应该进入您的连接 URL。表格不能包含点。该格式通常为[database].[table]
  • 你是对的,我意识到 bte3_myname 是架构,但没有它我的用户无法直接访问表..有没有办法在连接 url 中指定架构?我对 Oracle 数据库一无所知。我认为您可以为 postgres 数据库指定类似 &currentSchema 的内容
  • 您是否找到了有关您的 oracle 数据库驱动程序的 jdbc URI 格式的任何参考资料?
  • 不,还没有。我找到了这个站点:doc.nuodb.com/Latest/Content/…,它说我们可以这样做:jdbc://com.nuodb://host[:port]/database_name?[connection properties] --> jdbc:com.nuodb:/ /localhost/test?user=cloud&password=user&schema=mytest.我尝试在我的 url 连接“?schema=Bte3_utilitydata”的末尾添加它,但我得到了一个 SID 无效错误
  • 1) 用户和密码在 Connect 配置中是独立的属性。不要将它们以明文形式放在 url 中。 2)你缺少一个端口号 3)我想你想要这个jdbc:com.nuodb://localhost/Bte3_myname?schema=centrallogging

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


【解决方案1】:

如果 bte3_myname 实际上是您的架构,这可能会起作用

"table.name.format": "bte3_myname.${topic}"

(加或加一个下划线)。

我还注意到您使用的是混合大小写 - 因此您可能需要相应地设置“quote.sql.identifiers”。

【讨论】:

    【解决方案2】:

    如果主题名称由于命名原因包含点,但表名只是它的一部分,如topic.prefix.MY_TABLE_NAME.topic.suffix,则可以配置接收器连接器使用RegexRouter 转换,可以提取@987654323 @ 用于下沉操作。

    转换可能如下所示:

    "transforms": "changeTopicName",
    "transforms.changeTopicName.type": org.apache.kafka.connect.transforms.RegexRouter",
    "transforms.changeTopicName.regex": "topic.prefix.(MY_.*).topic.suffix",
    "transforms.changeTopicName.replacement": "$1",
    

    那么连接器将使用MY_TABLE_NAME 作为表名。

    附:确实,应该更聪明地定义正则表达式,但这取决于情况,对吧? ;)

    【讨论】:

      猜你喜欢
      • 2019-06-17
      • 2020-01-11
      • 2020-08-08
      • 2019-11-15
      • 2018-02-06
      • 2019-10-15
      • 2018-05-01
      • 2021-03-04
      • 2021-05-07
      相关资源
      最近更新 更多