【问题标题】:Can I define schema on kafka topic?我可以在 kafka 主题上定义模式吗?
【发布时间】:2020-07-27 00:08:20
【问题描述】:

我像这样向 kafka 主题发送数据和值模式:

./bin/kafka-avro-console-producer \
  --broker-list 10.0.0.0:9092 --topic orders \
  --property parse.key="true" \
  --property key.schema='{"type":"record","name":"key_schema","fields":[{"name":"id","type":"int"}]}' \
  --property key.separator="$" \
  --property value.schema='{"type":"record","name":"myrecord","fields":[{"name":"id","type":["null","int"],"default": null},{"name":"product","type": ["null","string"],"default": null}, {"name":"quantity", "type":  ["null","int"],"default": null}, {"name":"price","type":  ["null","int"],"default": null}]}' \
  --property schema.registry.url=http://10.0.0.0:8081

然后我从 kafka 获取此接收器属性的数据:

{
  "name": "jdbc-oracle",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics": "orders",
    "connection.url": "jdbc:oracle:thin:@10.1.2.3:1071/orac",
    "connection.user": "[redact]",
    "connection.password": "[redact]",
    "auto.create": "true",
    "delete.enabled": "true",
    "pk.mode": "record_key",
    "pk.fields": "id",
    "insert.mode": "upsert",
    "name": "jdbc-oracle"
  },
  "tasks": [
    {
      "connector": "jdbc-oracle",
      "task": 0
    }
  ],
  "type": "sink"
}

但我想从没有 value.schema 的 kafka 获取 json。如果我将 kafka 主题放在这个 json 数据中

{"id":9}${"id": {"int":9}, "product": {"string":"Yağız Gülbahar"}, "quantity": {"int":1071}, "price": {"int":61}}

如何从 kafka 获取这些数据并将 oracle 与 confluent jdbc sink 一起使用。

我想在 Kafka Connect 端制作架构?

另一件事是我可以从一个 kafka 主题中获取两种不同类型的数据,并且它可以通过 jdbc sink 在 oracle 端进入两个不同的表。

【问题讨论】:

  • 您的意思是,如何在不提供有效负载架构的情况下生成 Avro 消息?
  • plugin.path 不属于连接器配置

标签: apache-kafka apache-kafka-connect confluent-schema-registry


【解决方案1】:

如果您的源主题包含未声明架构的 JSON 数据,则必须添加该架构,然后才能使用 JDBC Sink。

选项包括:

  1. ksqlDB,如下图:https://www.youtube.com/watch?v=b-3qN_tlYR4&t=981s
  2. Kafka Connect 的单消息转换 功能。没有 SMT 附带 Apache Kafka 可以执行此操作,但有 prototypes out there 可以执行此操作。
  3. 其他流处理,例如卡夫卡流

编辑

我的意思是我可以从一个 kafka 主题定义两个不同的 jdbc sink 到不同的 oracle 表

是的,每个主题都可以被多个接收器消费。 table.name.format 配置选项可用于根据需要将主题路由到不同的表名。

【讨论】:

  • 感谢您的回答我明白必须有关于融合 jdbc 接收器主题的架构?并且融合的 jdbc 必须从一个主题进行一次转换。我说的对吗
  • "confluent jdbc 必须从一个主题进行一次转换。"你能更清楚地解释你在问什么吗?
  • “必须有关于融合 jdbc 接收器的主题模式?”正确。
  • 我已经更新了我的答案。在 StackOverflow 上,如果您对原始问题有新问题,您真的应该发布一个新问题 :)
  • 感谢您的回答,还有一件事可能我错了,我只是使用 jdbc-sink 将数据从 kafka 传输到 oracle(带有 jdbc 的任何数据库),我在写吗? ksql 或 ksqldb 是不同的东西?
猜你喜欢
  • 2018-11-07
  • 2016-01-02
  • 2020-08-06
  • 1970-01-01
  • 1970-01-01
  • 2021-01-28
  • 2020-09-11
  • 1970-01-01
  • 2022-06-16
相关资源
最近更新 更多