Athena与EMR PySpark查询S3数据结果不一致,PySpark返回全Null
排查PySpark读取S3 Parquet返回全Null但Athena正常的问题
我来帮你梳理几个关键的排查方向,结合你给出的信息,这些点大概率能定位到问题:
1. 先检查列名与Schema的匹配性
首先注意到你在Athena的查询里两个字段都用了<real_col1>,但PySpark脚本里换成了<real_col1>和<real_col2>——这会不会是笔误?如果<real_col2>在实际Parquet数据里本身无值,或者列名拼写/大小写不对,自然会返回Null。
建议先执行df1.printSchema()打印Spark读取Parquet后的Schema,对比Athena表的结构(用Athena的DESCRIBE <my Athena table name>查看):
- 确认列名大小写完全一致:比如Athena里列名是
Real_Col1,但Spark读出来是real_col1,而你查询用的是大写形式,就会匹配不到数据。 - 确认列的数据类型匹配:比如Athena里是STRING类型,但Spark读出来是INT,也可能导致读取异常返回Null。
另外你设置了spark.sql.hive.caseSensitiveInferenceMode=NEVER_INFER,这个参数会关闭Spark对大小写不敏感的推断,所以列名必须完全匹配才能读到数据。
2. 验证分区过滤的有效性
你用partition_datetime LIKE '2019-01-01-14'做过滤,但Spark直接读取Parquet时,分区列的处理逻辑和Athena有差异:
- 先确认S3路径是否包含分区目录结构,比如是否是
s3://path/partition_datetime=2019-01-01-14/?如果你的<s3:path to my data>只到父目录,要确认Spark是否自动识别了分区列。 - 可以先去掉过滤条件,执行
df1.select("partition_datetime").distinct().show(),看看实际的分区值格式,是不是和你写的2019-01-01-14完全一致。 - 也可以试试改用等值过滤:
WHERE partition_datetime = '2019-01-01-14',避免LIKE可能带来的匹配问题。
3. 排查Parquet文件的兼容性问题
Athena和EMR Spark使用的Parquet解析库版本可能存在差异,导致Spark无法正确解析某些Parquet元数据:
- 试试在读取Parquet时加上
mergeSchema参数:
这个参数会合并所有Parquet文件的Schema,避免因为部分文件Schema缺失导致的Null。df1 = spark.read.parquet("<s3:path to my data>", mergeSchema=True) - 你设置了
spark.sql.hive.convertMetastoreParquet=false,这个参数会让Spark使用Hive的Parquet SerDe读取数据,而非Spark原生解析器。可以临时注释掉这个参数,重新运行脚本看看是否正常。 - 还可以禁用向量式读取试试:
df1 = spark.read.option("spark.sql.parquet.enableVectorizedReader", "false").parquet("<s3:path to my data>")
4. 确认数据路径的准确性
确保PySpark里的<s3:path to my data>和Athena表的存储路径完全一致:
- 在Athena里执行
SHOW CREATE TABLE <my Athena table name>,查看表的LOCATION,对比你PySpark里的路径是否完全相同(S3路径是大小写敏感的)。 - 如果路径正确,先执行
df1.count()看看总数据量,如果count为0,说明根本没读到数据;如果count有数值但列都是Null,那就是列的匹配问题。
5. 检查Hive元数据与直接读Parquet的差异
Athena基于Hive元数据查询,而你PySpark是直接读取Parquet文件,两者处理逻辑有区别:
- Athena的表可能定义了自定义SerDe或类型转换规则,比如把Parquet里的二进制类型转成STRING,而Spark直接读的时候没有做这个转换,导致显示为Null。可以对比Athena表的DDL和Spark的Schema,看看有没有这类差异。
内容的提问来源于stack exchange,提问作者Thom Rogers
相关产品推荐
相关产品推荐

