如何在AWS Athena的OpenCSVSerde中跳过类型不匹配的行?
Athena读取CSV时跳过类型不匹配行的解决方法
问题描述
我在Athena中创建了一张从S3文件夹内的gzip压缩CSV文件读取数据的外部表,建表语句如下:
CREATE external TABLE IF NOT EXISTS `mydatabase`.`mytable` ( `messageId` string, `sourceCategory` string, `messageTime` string, `_messagetimepoch` string, `actallocmib` float, `activity` string, `bottom` integer, `numexecutions` integer, `numnodes` integer, `numprofilednodes` integer, `random` string, `tenantid` string, `usertime` integer, `version` string, `workunitduration` integer ) ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.OpenCSVSerde' WITH SERDEPROPERTIES ( 'separatorChar' = ',', 'quoteChar' = '\"', 'escapeChar' = '\\', 'serialization.format' = ',' ) STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat' LOCATION 's3:location' TBLPROPERTIES ( 'has_encrypted_data'='false', 'compression.type'='gzip', 'serialization.null.format'='', 'skip.header.line.count'='1' );
该表原本运行正常,直到出现以下错误:
HIVE_BAD_DATA: Error reading field value: Cannot convert value "hostPid":9912} of type String to a REAL value
我尝试在SERDEPROPERTIES中添加"ignore.malformed.json" = "true",但并未生效。需要实现当某列值与对应列类型不匹配时跳过整行。
解决方案
1. 理解无效参数的原因
ignore.malformed.json是JSON SerDe专属的配置参数,对处理CSV的OpenCSVSerDe完全不起作用,所以添加这个参数无法解决问题。
2. 使用Athena表属性跳过坏行
在表的TBLPROPERTIES中添加'skip.bad.record'='true',这个属性会让Athena在遇到解析或类型转换错误时自动跳过对应的行。
如果是新建表,修改后的TBLPROPERTIES部分如下:
TBLPROPERTIES ( 'has_encrypted_data'='false', 'compression.type'='gzip', 'serialization.null.format'='', 'skip.header.line.count'='1', 'skip.bad.record'='true' -- 新增该属性 );
如果是已存在的表,可通过ALTER语句修改:
ALTER TABLE `mydatabase`.`mytable` SET TBLPROPERTIES ('skip.bad.record'='true');
3. 稳妥替代方案:先读为字符串再转换过滤
如果上述方法仍不生效,可以采用更灵活的方式:先将所有列定义为STRING类型,然后在查询时使用TRY_CAST函数尝试转换为目标类型,过滤掉转换失败的行。
步骤1:创建全字符串类型的表
CREATE external TABLE IF NOT EXISTS `mydatabase`.`mytable_string` ( `messageId` string, `sourceCategory` string, `messageTime` string, `_messagetimepoch` string, `actallocmib` string, `activity` string, `bottom` string, `numexecutions` string, `numnodes` string, `numprofilednodes` string, `random` string, `tenantid` string, `usertime` string, `version` string, `workunitduration` string ) ROW FORMAT SERDE 'org.apache.hadoop.hive.serde2.OpenCSVSerde' WITH SERDEPROPERTIES ( 'separatorChar' = ',', 'quoteChar' = '\"', 'escapeChar' = '\\', 'serialization.format' = ',' ) STORED AS INPUTFORMAT 'org.apache.hadoop.mapred.TextInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat' LOCATION 's3:location' TBLPROPERTIES ( 'has_encrypted_data'='false', 'compression.type'='gzip', 'serialization.null.format'='', 'skip.header.line.count'='1' );
步骤2:查询时转换并过滤
SELECT messageId, sourceCategory, messageTime, _messagetimepoch, TRY_CAST(actallocmib AS float) AS actallocmib, activity, TRY_CAST(bottom AS integer) AS bottom, TRY_CAST(numexecutions AS integer) AS numexecutions, TRY_CAST(numnodes AS integer) AS numnodes, TRY_CAST(numprofilednodes AS integer) AS numprofilednodes, random, tenantid, TRY_CAST(usertime AS integer) AS usertime, version, TRY_CAST(workunitduration AS integer) AS workunitduration FROM mydatabase.mytable_string -- 过滤所有转换失败的行 WHERE TRY_CAST(actallocmib AS float) IS NOT NULL AND TRY_CAST(bottom AS integer) IS NOT NULL AND TRY_CAST(numexecutions AS integer) IS NOT NULL AND TRY_CAST(numnodes AS integer) IS NOT NULL AND TRY_CAST(numprofilednodes AS integer) IS NOT NULL AND TRY_CAST(usertime AS integer) IS NOT NULL AND TRY_CAST(workunitduration AS integer) IS NOT NULL;
TRY_CAST函数在转换失败时会返回NULL,通过WHERE子句过滤掉包含NULL的行,即可实现跳过类型不匹配的原始数据行。
内容的提问来源于stack exchange,提问作者Adem Youssef
相关产品推荐
相关产品推荐

