【发布时间】: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