【问题标题】:Kafka: All messages failing in stream while data in topicsKafka:所有消息在流中失败,而主题中的数据
【发布时间】:2020-07-14 12:59:34
【问题描述】:

我有主题 post_users_t 并且在使用 PRINT 命令时我得到了

rowtime: 4/2/20 2:03:48 PM UTC, key: <null>, value: {"userid": 6, "id": 8, "title": "testest", "body": "Testingmoreand more"}
rowtime: 4/2/20 2:03:48 PM UTC, key: <null>, value: {"userid": 7, "id": 11, "title": "testest", "body": "Testingmoreand more"}

然后我用这个创建一个流:

CREATE STREAM userstream (userid INT, id INT, title VARCHAR, body VARCHAR)
    WITH (KAFKA_TOPIC='post_users_t',
          VALUE_FORMAT='JSON');

但我无法从中选择任何内容,当我 DESCRIBE EXTENDED 它时,所有消息都失败了。

consumer-messages-per-sec:      1.06 consumer-total-bytes:    116643 consumer-total-messages:      3417     last-message: 2020-04-02T14:08:08.546Z
consumer-failed-messages:      3417 consumer-failed-messages-per-sec:      1.06      last-failed: 2020-04-02T14:08:08.56Z

我在这里做错了什么?

下面有额外的信息!

从头开始打印主题:

ksql> print 'post_users_t' from beginning limit 2;
Key format: SESSION(AVRO) or HOPPING(AVRO) or TUMBLING(AVRO) or AVRO or SESSION(PROTOBUF) or HOPPING(PROTOBUF) or TUMBLING(PROTOBUF) or PROTOBUF or SESSION(JSON) or HOPPING(JSON) or TUMBLING(JSON) or JSON or SESSION(JSON_SR) or HOPPING(JSON_SR) or TUMBLING(JSON_SR) or JSON_SR or SESSION(KAFKA_INT) or HOPPING(KAFKA_INT) or TUMBLING(KAFKA_INT) or KAFKA_INT or SESSION(KAFKA_BIGINT) or HOPPING(KAFKA_BIGINT) or TUMBLING(KAFKA_BIGINT) or KAFKA_BIGINT or SESSION(KAFKA_DOUBLE) or HOPPING(KAFKA_DOUBLE) or TUMBLING(KAFKA_DOUBLE) or KAFKA_DOUBLE or SESSION(KAFKA_STRING) or HOPPING(KAFKA_STRING) or TUMBLING(KAFKA_STRING) or KAFKA_STRING
Value format: AVRO or KAFKA_STRING
rowtime: 4/2/20 1:04:08 PM UTC, key: <null>, value: {"userid": 1, "id": 1, "title": "loremit", "body": "loremit heiluu ja paukkuu"}
rowtime: 4/2/20 1:04:08 PM UTC, key: <null>, value: {"userid": 2, "id": 2, "title": "lorbe", "body": "larboloilllaaa"}

【问题讨论】:

  • 你能把PRINT post_users_t FROM BEGINNING LIMIT 2的完整输出贴出来吗
  • 另外你运行的是什么版本的 Confluent 平台?
  • 嘿@RobinMoffatt。再次感谢您的快速答复。我从一开始就添加了打印,实际上我之前没有使用过限制,所以我之前没有看到过这个输出。
  • 哦,是的,我在 docker 上运行一切,你可以在这里找到 docker-compose:github.com/Itzblend/KafkaPOC/blob/master/docker-compose.yml

标签: apache-kafka apache-kafka-connect confluent-platform ksqldb


【解决方案1】:

根据 ksqlDB 对该主题的检查输出,您的数据在 Avro 中被序列化:

Value format: AVRO or KAFKA_STRING

但是您已经创建了 STREAM 并指定了 VALUE_FORMAT='JSON'。这将导致反序列化错误,如果您运行docker-compose logs -f ksqldb-server,您将在尝试查询流时​​看到被写出。

由于您使用的是 Avro,因此无需指定架构。试试这个:

CREATE STREAM userstream 
   WITH (KAFKA_TOPIC='post_users_t',
         VALUE_FORMAT='AVRO');

【讨论】:

  • 太棒了!让它完全按照你的建议工作。我想知道是什么将它指定为 AVRO。再次感谢您的快速回答,我想说我喜欢您的演讲。他们是所有技术领域中最好的之一!
猜你喜欢
  • 2023-01-07
  • 2021-03-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-08-16
  • 2018-10-27
  • 1970-01-01
  • 2016-01-07
相关资源
最近更新 更多