【问题标题】:Apache Flink : Handle bad avro records in confluent-avro from KafkaApache Flink:处理来自 Kafka 的 confluent-avro 中的坏 avro 记录
【发布时间】:2021-11-30 20:12:44
【问题描述】:

我已经使用 Flink 的表 API 创建了一个表。

CREATE TABLE recommendations (
 ...

) WITH (
  'connector' = 'kafka',
  'topic' = 'my_kafka_topic',
  'properties.bootstrap.servers' = 'localhost:9092',
  'properties.group.id' = 'testGroup',
  'properties.security.protocol' = 'SASL_PLAINTEXT',
  'properties.sasl.kerberos.service.name' = 'kafka',
  'scan.startup.mode' = 'latest-offset',
  'value.format' = 'avro-confluent',
  'value.avro-confluent.url' = 'http://schema-registry-address',
  'value.fields-include' = 'EXCEPT_KEY'
);

当运行 SQL 来查看记录时,我得到:

Flink SQL> select * from default_catalog.default_database.recommendations ;
[ERROR] Could not execute SQL statement. Reason:
java.lang.ArrayIndexOutOfBoundsException: -25

Flink SQL> select * from default_catalog.default_database.recommendations ;
[ERROR] Could not execute SQL statement. Reason:
java.io.IOException: Failed to deserialize Avro record.

我知道有一些 BAD avro 记录被推送到 Kafka 主题中。在 JSON 格式中,可以通过设置来跳过/过滤这些记录 'json.ignore-parse-errors' = 'true'从 confluent-avro 格式读取时,我们有什么办法可以跳过这些记录?

这并不理想,但不幸的是,尽管有架构注册表,但我无法控制推送到 Kafka 的内容。

【问题讨论】:

    标签: apache-kafka apache-flink avro flink-streaming flink-sql


    【解决方案1】:

    AVRO 目前没有这样的选项。在https://issues.apache.org/jira/browse/FLINK-20091有一张公开的票。

    【讨论】:

      猜你喜欢
      • 2017-05-06
      • 2016-11-10
      • 2017-02-17
      • 2021-10-23
      • 2021-05-11
      • 2015-05-01
      • 2021-06-14
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多