【问题标题】:Consumer_failed_message in kafka stream: Records not pushed from topickafka 流中的 Consumer_failed_message:记录未从主题推送
【发布时间】:2020-06-03 09:38:21
【问题描述】:

我有一个从 IBM 大型机 IIDR 向 Kafka 主题发送记录的流程。进入 Kafka 主题的消息的value_format 是 AVRO,密钥也是 AVRO 格式。记录被推送到 Kafka 主题中。我有一个与该主题相关的流。但是记录不会传递到流中。 test_iidr 主题示例 -

rowtime: 5/30/20 7:06:34 PM UTC, key: {"col1": "A", "col2": 1}, value: {"col1": "A", "col2": 11, "col3": 2, "iidr_tran_type": "QQ", "iidr_a_ccid": "0", "iidr_a_user": " ", "iidr_src_upd_ts": "2020-05-30 07:06:33.262931000", "iidr_a_member": " "}

流中的 value_format 是 AVRO 并且列名都被检查。

流创建查询 -

CREATE STREAM test_iidr (
  col1 STRING, 
  col2 DECIMAL(2,0),
  col3 DECIMAL(1,0),
  iidr_tran_type STRING,
  iidr_a_ccid STRING,
  iidr_a_user STRING,
  iidr_src_upd_ts STRING,
  iidr_a_member STRING) 
WITH (KAFKA_TOPIC='test_iidr', PARTITIONS=1, REPLICAS=3, VALUE_FORMAT='AVRO');

由于KEY 未在WITH 语句中提及,因此无法从主题加载到流中? 模式注册表中注册了 test_iidr-value 和 test_iidr-key 主题。

Kafka-connect 泊坞窗中的 key.converter 和 value.converter 设置为 - org.apache.kafka.connect.json.JsonConverter。这是JsonConverter 创建这个问题吗?

我用不同的流创建了一个完全不同的管道,并使用insert into 语句手动插入了相同的数据。有效。只有 IIDR 流不起作用,并且记录没有从主题推送到流中。

我正在使用 Confluent kafka 5.5.0 版。

【问题讨论】:

  • 问题已解决。看起来由于 DECIMAL 和 INT 存在反序列化错误。源正在发送 INT 值,我们将 DECIMAL 作为数据类型。

标签: ksqldb confluent-platform kafka-topic


【解决方案1】:

连接配置中的JsonConverter 很可能会将您的 Avro 数据转换为 JSON。

要确定键和值序列化格式,您可以使用PRINT 命令(我可以看到您已经运行了该命令)。 PRINT 运行时会输出键值格式。例如:

ksql> PRINT some_topic FROM BEGINNING LIMIT 1;
Key format: JSON or KAFKA_STRING
Value format: JSON or KAFKA_STRING
rowtime: 5/30/20 7:06:34 PM UTC, key: {"col1": "A", "col2": 1}, value: {"col1": "A", "col2": 11, "col3": 2, "iidr_tran_type": "QQ", "iidr_a_ccid": "0", "iidr_a_user": " ", "iidr_src_upd_ts": "2020-05-30 07:06:33.262931000", "iidr_a_member": " "}

所以首先要检查的是通过 PRINT 输出键和值的格式,然后相应地更新您的 CREATE 语句。

注意,ksqlDB尚不支持 Avro/Json 键,因此您可能希望/需要重新分区数据,请参阅:https://docs.ksqldb.io/en/latest/developer-guide/syntax-reference/#what-to-do-if-your-key-is-not-set-or-is-in-a-different-format

旁注:如果值的架构存储在架构注册表中,那么您不需要在 CREATE 语句中定义列,因为 ksqlDB 将从架构注册表加载列

旁注:对于现有主题,您不需要在 WITH 子句中使用 PARTITIONS=1, REPLICAS=3,只有当您希望 ksqlDB 为您创建主题时。

【讨论】:

  • 感谢您的建议。我做了改变。问题在于数据类型 DECIMAL。源作为 INTEGER 而不是 DECIMAL 发送,并且 Kafka 流在模式中有 DECIMAL。我们将其更新为 INT。成功了
  • 好东西。真高兴你做到了。你能把问题标记为已回答吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-12-19
  • 1970-01-01
  • 2018-09-09
  • 1970-01-01
  • 2019-03-27
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多