近日,多位数据工程师在技术社区报告了一个令人困惑的Parquet格式问题:当使用嵌套JSON字段(即record type)写入Parquet文件时,尽管在查询测试结果中这些字段能够正确填充,但最终输出的Parquet文件中相关字段却全部显示为null。该问题导致下游的数据分析、机器学习模型训练等任务无法获取真实数据,引发广泛关注。
问题重现:从正常到异常的“隐形”丢失
据悉,问题最初在数仓迁移项目中暴露。工程师将业务系统产生的JSON嵌套数据通过ETL管道写入Parquet格式存储,并在中间环节使用Spark SQL进行数据验证——查询结果中,嵌套字段“user.address.city”等均正常显示具体值。然而,当直接读取同一Parquet文件时,却发现这些record type字段全部为空(null),而简单类型字段(如string、int)则完全正常。
进一步测试发现,该问题在Hive、Presto、Trino等多个查询引擎中均可复现,且与数据量大小、嵌套层级无关。只要字段被定义为struct或嵌套的record type,Parquet输出时就会“漏掉”这些值,但同一条记录在内存或临时表(如内存表、Avro格式文件)中却表现正常。
根源追踪:Parquet投影与类型推导的冲突
多位开源贡献者将矛头指向Parquet写入时的schema推导机制。Parquet作为列式存储格式,在写入嵌套类型时依赖严格的元数据定义(如Repetition和Definition levels)。当数据源(如JSON)本身具有动态或松散的结构,且ETL工具采用自动schema推断时,可能出现以下情况:
- 写入器将嵌套字段错误地标记为optional但未处理实际值;
- 读取器在解析时因类型不匹配(例如JSON中字符串被推测为struct)而回退为null;
- Parquet的投影缩减(projection pushdown)优化在嵌套字段上产生副作用,导致中间层丢弃有效数据。
此外,某些引擎(如Spark 3.x)的Parquet写入器在启用“vectorized reader”或“adaptive query execution”时,会预先根据统计信息推断字段是否为空——若样本中少量行无值,可能误判整个字段为“全空”,进而压缩为null。
影响范围:从报表异常到模型灾难
该问题已在多家企业的生产环境中引发数据一致性事故。某电商平台数据团队反馈,其用户行为分析报表中所有“浏览记录”嵌套字段突然归零,导致次日转化率指标骤降20%,经排查才确认是Parquet输出错误。另一家金融科技公司则发现,其风控模型输入的特征向量中,约60%的嵌套特征变为null,模型失效三天后才被修复。
更棘手的是,由于查询测试阶段结果正常,工程师往往第一时间怀疑业务代码问题,而非存储格式故障,平均定位时间长达6-8小时。部分团队不得不回滚至JSON Lines格式,牺牲查询性能以换取数据正确性。
临时解决方案与长期改进
截至目前,主要数据处理框架已发布或建议以下缓解措施:
- 禁用自动schema推断:在Spark/Presto中明确指定Parquet写入的schema,而非依赖JSON自推导。例如使用
struct<address:struct<city:string,street:string>>预定义类型。 - 转换嵌套为扁平结构:将嵌套字段通过JSON函数或UDF展开为多个独立列,但会增加存储空间。
- 回退至Avro或ORC格式:Avro在嵌套类型处理上更稳定,但压缩率和查询性能略逊于Parquet。
- 升级引擎版本:Apache Spark 3.5.1、Trino 436等已修复部分相关bug,建议用户检查发行说明。
长期来看,Apache Parquet社区正在推进其“Logical Type”的强化,计划增加对JSON嵌套的显式支持(如MAP<STRING, STRUCT<...>>),避免隐式类型转换导致的null问题。同时,各大云厂商(如AWS Glue、Databricks)也在其Parquet写入SDK中增加了字段级校验日志,便于快速定位。
专家建议:测试不可止于中间环节
“这个案例再次提醒我们,大数据管道的验证不能仅依赖中间结果查询。”某资深数据架构师在技术博客中写道,“最终文件的格式一致性测试应纳入CI/CD流程,尤其是使用非范式数据源(如JSON、XML)时,必须对Parquet文件做实际读取校验,包括schema对比和随机抽样值比对。”
目前,该问题已被标记为Apache Spark JIRA(SPARK-43258)及Trino GitHub Issue(#18502)的高优先级议题,预计未来1-2个月内会有官方修复版本发布。对于正在遭受数据丢失困扰的团队,建议立即采用手动指定schema的写入策略,并在生产环境中对Parquet输出文件进行全量校验,避免连锁故障。