【发布时间】:2022-02-05 06:32:16
【问题描述】:
当我将 kafka JDBC 连接器运行到 PSQL 时出现以下错误:
带有 schemas.enable 的 JsonConverter 需要“schema”和“payload” 字段,并且可能不包含其他字段。如果你想 反序列化纯 JSON 数据,设置 schemas.enable=false 在您的 转换器配置。
但是我的主题包含以下消息结构,其中添加了一个模式,就像它在网上展示的一样:
rowtime: 2022/02/04 12:45:48.520 Z, key: , value: "{"schema": {“类型”:“结构”,“字段”:[{“类型”:“int”,“字段”: “ID”、“可选”:false}、{“类型”:“日期”、“字段”: “日期”,“可选”:假},{“类型”:“varchar”, “字段”:“ICD”,“可选”:假},{“类型”:“int”, “字段”:“CPT”,“可选”:假},{“类型”:“双”, “字段”:“成本”,“可选”:假}],“可选”:假, "name": "test"}, "payload": {"ID": "24427934", “日期”:“2019-05-22”,“ICD”:“883.436”,“CPT”: "60502", "成本": "1374.36"}}", 分区:0
我对连接器的配置是:
curl -X PUT http://localhost:8083/connectors/claim_test/config \
-H "Content-Type: application/json" \
-d '{
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"connection.url":"jdbc:postgresql://localhost:5432/ae2772",
"key.converter":"org.apache.kafka.connect.json.JsonConverter",
"value.converter":"org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable":"true",
"topics":"test_7",
"auto.create":"true",
"insert.mode":"insert"
}'
经过一些更改后,我现在收到以下消息:
WorkerSinkTask{id=claim_test} Error converting message value in topic 'test_9' partition 0 at offset 0 and timestamp 1644005137197: Unknown schema type: int
【问题讨论】:
-
您是否检查了所有您主题中的消息都遵循这种格式?
-
是的 - 我已经删除了以前的主题,当前主题有一个生产者从 python 发布 json 消息。消息的结构都是一致的。
-
我知道收到以下消息:WorkerSinkTask{id=claim_test} 在偏移量 0 和时间戳 1644005137197 的主题“test_9”分区 0 中转换消息值时出错:未知模式类型:int - 这里的任何帮助都是很棒
标签: apache-kafka apache-kafka-connect jdbc-postgres