【发布时间】:2022-01-16 19:38:19
【问题描述】:
尝试在 ksql CLI 上创建 ElasticSearch 接收器连接器时,我收到以下错误:
ERROR WorkerSinkTask{id=SINK_ELASTIC_TEST_JSON_A-0} 转换错误 主题“REROUTES_TABLE”分区 0 中偏移量 939 处的消息值和 时间戳 1641056495920:将 byte[] 转换为 Kafka Connect 数据 由于主题 REROUTES_TABLE 的序列化错误而失败: (org.apache.kafka.connect.runtime.WorkerSinkTask)
引起:org.apache.kafka.common.errors.SerializationException: 为 id 30 反序列化 JSON 消息时出错原因: java.net.ConnectException:连接被拒绝(连接被拒绝)在 java.base/java.net.PlainSocketImpl.socketConnect(Native Method) 在 java.base/java.net.AbstractPlainSocketImpl.doConnect(AbstractPlainSocketImpl.java:399)
它的创建命令如下所示:
CREATE SINK CONNECTOR SINK_ELASTIC_TEST_JSON_A WITH (
'connector.class' = 'io.confluent.connect.elasticsearch.ElasticsearchSinkConnector',
'connection.url' = 'http://elasticsearch:9200',
'key.converter' = 'org.apache.kafka.connect.storage.StringConverter',
'value.converter' = 'io.confluent.connect.json.JsonSchemaConverter',
'value.converter.schema.registry.url' = 'http://localhost:8081',
'value.converter.schemas.enable' = 'true',
'type.name' = '_doc',
'topics' = 'REROUTES_TABLE',
'key.ignore' = 'false',
'schema.ignore' = 'false'
);
数据看起来像这样:
ksql> print REROUTES_TABLE from beginning limit 1;
密钥格式:KAFKA_INT 或 KAFKA_STRING 值格式:JSON_SR 或 KAFKA_STRING 行时间:2021/12/26 06:22:33.726 Z,键:0,值:{“STEP_CNT”:1,“TOT_LEN”:0.0013573977968634994},分区:0 主题打印停止
主题值的架构是:
{"subject":"REROUTES_TABLE-value","version":1,"id":30,"schemaType":"JSON","schema":"{"type":"object","properties ":{"STEP_CNT":{"connect.index":0,"oneOf":[{"type":"null"},{"type":"integer","connect.type":"int64"} ]},"TOT_LEN":{"connect.index":1,"oneOf":[{"type":"null"},{"type":"number","connect.type":"float64"} ]}}}"}
REROUTES_TABLE 建立在流上,对流数据进行了一些聚合。
我有点怀疑反序列化器无法理解一个空值,但由于 REROUTES_TABLE 能够在流上执行聚合,空值是如何以及从哪里来的,更重要的是如何解决这个问题(甚至如果我对 null 的假设不正确)?
【问题讨论】:
-
我会尝试使用
kcat -p 0 -o 939检查错误中报告的失败偏移量的确切记录是什么。您是否尝试过从主题的开头开始,主题中是否有 938 条有效记录,或者这是您遇到的第一个错误? -
是的,我确实尝试从主题的开头开始,但发生了同样的错误。另外,我错过了在日志上报告一些信息,我还看到:Caused by: java.net.ConnectException: Connection denied (Connection denied)(也更新了描述)不确定它无法连接到什么。跨度>
-
1) 您可以删除
value.converter.schemas.enable2) 错误表明value.converter.schema.registry.url指向您的架构注册表的错误地址,或者它没有在您提供的地址上运行 -
哦,谢谢你的指点!!通过在 value.converter.schema.registry.url 中用 schema-registry 替换 localhost 解决了这个问题。不过这只是侥幸,我想知道它们之间的区别是什么。
-
有什么区别?
标签: apache-kafka apache-kafka-connect confluent-schema-registry