如何获取FlinkSQL中被value.json.ignore-parse-errors忽略的错误数据日志?
获取Flink SQL中被忽略的JSON解析错误记录日志
1. 开启Flink组件的错误日志输出
直接调整Flink的日志配置文件(log4j.properties或log4j2.properties),添加以下配置:
# 开启JSON格式解析模块的日志,级别设为INFO即可捕获错误记录 log4j.logger.org.apache.flink.formats.json=INFO log4j.logger.org.apache.flink.table.runtime.typeutils=INFO
重启任务后,解析错误时的原始数据、错误详情都会被输出到Flink的任务日志中,直接在日志文件里搜索相关关键词就能定位问题数据。
2. 用侧输出流精准收集错误数据
如果不想依赖日志,想把错误数据单独存储以便反馈给上游,可以用侧输出流的方式:
- 先创建一张专门存储错误记录的表(示例为写入Kafka):
CREATE TABLE error_records ( raw_data STRING, error_msg STRING, capture_time TIMESTAMP(3) ) WITH ( 'connector' = 'kafka', 'topic' = 'flink_json_parse_errors', 'properties.bootstrap.servers' = '你的Kafka地址', 'format' = 'json' );
- 然后在主查询中用
TRY_CAST判断解析是否成功,将错误数据路由到错误表:
INSERT INTO main_table, error_records SELECT TRY_CAST(date_field AS DATE) AS valid_date, -- 映射其他正常字段 NULL AS raw_data, NULL AS error_msg, NULL AS capture_time FROM source_table WHERE TRY_CAST(date_field AS DATE) IS NOT NULL UNION ALL SELECT NULL AS valid_date, -- 其他字段设为NULL TO_STRING(RAW()) AS raw_data, CONCAT('日期字段解析失败: ', CAST(ERROR() AS STRING)) AS error_msg, CURRENT_TIMESTAMP() AS capture_time FROM source_table WHERE TRY_CAST(date_field AS DATE) IS NULL;
错误数据会被单独写入指定的Kafka Topic,导出后就能直接给上游展示具体的错误内容。
3. 临时调整任务日志级别(无需修改全局配置)
如果不想改动全局日志配置,提交任务时可以通过参数临时生效:
./bin/flink run -d \ -Dlog4j.logger.org.apache.flink.formats.json=INFO \ 你的FlinkSQL作业.jar
内容的提问来源于stack exchange,提问作者Faisal Ahmed Siddiqui
相关产品推荐
相关产品推荐

