【问题标题】:Kafka jdbc sink connector to Postgres failing with "cross-database references are not implemented"到 Postgres 的 Kafka jdbc sink 连接器因“未实现跨数据库引用”而失败
【发布时间】:2022-01-07 21:06:33
【问题描述】:

我的设置基于 docker 容器 - 1 个 oracle db、kafka、kafka connect 和 postgres。我首先使用 oracle CDC 连接器来提供工作正常的 kafka。然后我试图阅读该主题并将其输入 Postgres。 当我启动连接器时,我得到:

"trace": "org.apache.kafka.connect.errors.ConnectException: 由于不可恢复的异常退出 WorkerSinkTask。\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:614 )\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:329)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:232)\ n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:201)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:185)\n\ tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:234)\n\tat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n\ tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java .base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.l ang.Thread.run(Thread.java:834)\n引起:org.apache.kafka.connect.errors.ConnectException:java.sql.SQLException:异常链:\norg.postgresql.util.PSQLException:错误:交叉数据库引用未实现:“ORCLCDB.C__MYUSER.EMP”\n 位置:14\n\n\tat io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:122)\n\tat org. apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:586)\n\t... 10 more\nCaused by: java.sql.SQLException: Exception chain:\norg.postgresql.util.PSQLException: ERROR : 跨数据库引用未实现:"ORCLCDB.C__MYUSER.EMP"\n 位置:14\n\n\tat io.confluent.connect.jdbc.sink.JdbcSinkTask.getAllMessagesException(JdbcSinkTask.java:150)\n\ tat io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:102)\n\t... 11 更多\n"

我的配置 json 看起来像:

{
"name": "SimplePostgresSink",
"config":{
  "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
  "name": "SimplePostgresSink",
  "tasks.max":1,
  "topics": "ORCLCDB.C__MYUSER.EMP",
  "key.converter": "io.confluent.connect.avro.AvroConverter",
  "key.converter.schema.registry.url": "http://schema-registry:8081",
  "value.converter": "io.confluent.connect.avro.AvroConverter",
  "value.converter.schema.registry.url": "http://schema-registry:8081",
  "confluent.topic.bootstrap.servers":"kafka:29092",
  "connection.url": "jdbc:postgresql://postgres:5432/postgres",
  "connection.user": "postgres",
  "connection.password": "postgres",
  "insert.mode": "upsert",
  "pk.mode": "record_value",
  "pk.fields": "I",
  "auto.create": "true",
  "auto.evolve": "true"
}

}

主题架构是:


{
  "type": "record",
  "name": "ConnectDefault",
  "namespace": "io.confluent.connect.avro",
  "fields": [
    {
      "name": "I",
      "type": {
        "type": "bytes",
        "scale": 0,
        "precision": 64,
        "connect.version": 1,
        "connect.parameters": {
          "scale": "0"
        },
        "connect.name": "org.apache.kafka.connect.data.Decimal",
        "logicalType": "decimal"
      }
    },
    {
      "name": "NAME",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "table",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "scn",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "op_type",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "op_ts",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "current_ts",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "row_id",
      "type": [
        "null",
        "string"
      ],
      "default": null
    },
    {
      "name": "username",
      "type": [
        "null",
        "string"
      ],
      "default": null
    }
  ]
}

我对 I 和 Name 列感兴趣

【问题讨论】:

  • 请显示主题中有哪些数据,以及您期望输出的 postgres 数据是什么

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


【解决方案1】:

查看日志我看到连接器在创建具有这样名称的表时出现问题:

[2022-01-07 23:56:11,737] 信息 JdbcDbWriter 已连接 (io.confluent.connect.jdbc.sink.JdbcDbWriter) [2022-01-07 23:56:11,759] 信息检查 PostgreSql 方言是否存在 TABLE "ORCLCDB"."C__MYUSER"."EMP" (io.confluent.connect.jdbc.dialect.GenericDatabaseDialect) [2022-01-07 23:56:11,764] 信息使用 PostgreSql 方言表“ORCLCDB”。“C__MYUSER”。“EMP”不存在(io.confluent.connect.jdbc.dialect.GenericDatabaseDialect) [2022-01-07 23:56:11,764] 信息使用 sql 创建表:CREATE TABLE "ORCLCDB"."C__MYUSER"."EMP" ( “我”十进制不为空, “名称”文本不为空, 主键("I","NAME")) (io.confluent.connect.jdbc.sink.DbStructure) [2022-01-07 23:56:11,765] WARN 创建失败,如果表已经存在,将尝试修改(io.confluent.connect.jdbc.sink.DbStructure) org.postgresql.util.PSQLException:错误:未实现跨数据库引用:“ORCLCDB.C__MYUSER.EMP”

【讨论】:

  • 正如目前所写,您的答案尚不清楚。请edit 添加其他详细信息,以帮助其他人了解这如何解决所提出的问题。你可以找到更多关于如何写好答案的信息in the help center
猜你喜欢
  • 2018-06-13
  • 2019-01-21
  • 2019-06-04
  • 1970-01-01
  • 2019-12-23
  • 2020-08-09
  • 2021-07-14
  • 2021-05-17
  • 1970-01-01
相关资源
最近更新 更多