【问题标题】:What is the simplest way to sync a Kafka KTable to a SQL database?将 Kafka KTable 同步到 SQL 数据库的最简单方法是什么?
【发布时间】:2019-01-05 04:12:49
【问题描述】:

我已经使用 KSQL 创建了一个流和一个聚合表。

  {  
   "ksql":"DROP Stream IF EXISTS StreamLegacyNames; DROP Stream IF EXISTS StreamLegacy; CREATE Stream  StreamLegacy (payload  STRUCT<AgeYr  varchar>)WITH (KAFKA_TOPIC='eip-legacy-13',VALUE_FORMAT='JSON' );  CREATE Stream  StreamLegacyNames As Select payload->AgeYr Age from StreamLegacy; Create Table DimAge As SELECT Age FROM StreamLegacyNames Group By Age;",
   "streamsProperties":{  
      "ksql.streams.auto.offset.reset":"earliest"
   }
}

将此代码导出到 sql 表的最简单方法是什么?我们正在为主题使用 jdbc 连接器,但我不清楚这是否适用于聚合的 KSQL 表(在此示例中为 DIMAGE)。

即使我在 jdbc 连接配置文件中将主题设置为 DIMAGE 和以下内容。

value.converter.schemas.enable=false

完整的配置文件是

connector.class=io.confluent.connect.jdbc.JdbcSinkConnector
connection.password=PASSWORD
auto.evolve=true
topics=DIMAGE
tasks.max=1
connection.user=USER
value.converter.schemas.enable=false
auto.create=true
connection.url=jdbc:sqlserver://SERVER

我在连接器中收到以下错误。

Caused by: org.apache.kafka.connect.errors.DataException: JsonConverter with schemas.enable requires "schema" and "payload" fields and may not contain additional fields. If you are trying to deserialize plain JSON data, set schemas.enable=false in your converter configuration.

通过 postman 的 KSQL 查询显示 KTABLE 的格式为

{"row":{"columns":["83"]},"errorMessage":null,"finalMessage":null}
{"row":{"columns":["74"]},"errorMessage":null,"finalMessage":null}
{"row":{"columns":["36"]},"errorMessage":null,"finalMessage":null}

【问题讨论】:

    标签: sql-server apache-kafka ksqldb


    【解决方案1】:

    当您在 KSQL 中 CREATE STREAM foo AS SELECT(“CSAS”)时,您正在创建一个新的 Kafka 主题并不断地使用 SELECT 语句的结果填充它。

    所以您只有一个 Kafka 主题,在您的情况下称为 STREAMLEGACYNAMES(KSQL 通常强制对象为大写)。您可以使用JDBC Sink connector 将此主题流式传输到目标 RDBMS,包括 MS SQL。

    【讨论】:

    • 嗨,谢谢。我想要完成的是使用 JDBC 接收器连接器将聚合的 KTABLE DimAge 发送到 sql。即使我设置 value.converter.schemas.enable=false,我也会收到一条错误消息,上面写着 Caused by: org.apache.kafka.connect.errors.DataException: JsonConverter with schemas.enable requires "schema" and "payload" fields并且可能不包含其他字段。如果您尝试反序列化纯 JSON 数据,请在转换器配置中设置 schemas.enable=false。
    • JDBC Sink 需要一个模式。如果您使用的是 KSQL,那么唯一的选择(也是一个不错的选择)是使用 Avro。如果您使用的是 JSON,那么 KSQL 生成的 JSON 不包括架构本身。更多详情请参阅本文confluent.io/blog/…
    【解决方案2】:

    说到底,KTable 只是另一个话题。您可以使用 KSQL PRINTkafka-console-consumer 查看 JDBC Sink 连接器将获取哪些数据。

    如果您假设 KSQL 表将与 SQL Server 表完全匹配,那么它不会。在 SQL Server 表中,您将拥有 KTable 上发生的每一个“事件行”,包括空值,因为 JDBC 接收器尚不支持删除。


    不确定您期望什么数据,但您可以做的是对您尝试捕获的事件执行窗口输出,然后您实际上是在将微批量插入到您的下游数据库中。

    【讨论】:

    • 嗨,谢谢我编辑了原始问题,以显示有关我正在尝试执行的操作和当前错误消息的更多详细信息。
    • 我假设您的密钥转换器也设置为 JSON,因此您还需要为其设置 schemas.enable
    猜你喜欢
    • 2011-08-08
    • 1970-01-01
    • 2012-07-04
    • 1970-01-01
    • 2011-05-13
    • 1970-01-01
    • 2011-01-15
    • 2020-05-13
    • 1970-01-01
    相关资源
    最近更新 更多