【发布时间】: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