You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 08:30:04