news 2026/9/11 20:51:05

Parquet转JSONL实战:Python生产级转换方案与踩坑指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Parquet转JSONL实战:Python生产级转换方案与踩坑指南

先说我最近遇到的一件事。数据平台跑完一批用户行为特征,产出的是一个接近 10GB 的 Parquet 文件,按天分区一共几十个小文件。下游业务方却说他们只需要 JSONL,因为他们要拿去做事件流回放,还要直接在命令行里用 jq 看明文。我当时的第一个反应是:这不就是把 Parquet 一行行读出来再写成 JSON 字符串吗?真做起来才发现,事情远没有“三行代码搞定”那么简单。好在这个需求非常典型——把 Parquet 转成 JSONL,在 Python 生态里就是一个“会者不难”的活,难的是你得知道每一步背后发生了什么,才知道什么时候该用哪种写法。

这篇文章就把我从头到尾踩过的坑、验证过的方案、以及最终沉淀下来的生产级写法完整讲一遍。不管你是刚入门 Python 的初学者,还是已经在数据管道里摸爬滚打了几年,只要手头有 Parquet 文件需要转成 JSONL,这篇文章的代码和经验都能直接拿去用。

1. 这个转换需求是怎么来的,以及两个格式的本质差异

先说人话:Parquet 是给机器看的高压缩列式存储,JSONL 是给人看的一行一个 JSON 对象的文本格式。两者没有谁替代谁,它们活在完全不同的工作流里。

我见过最常见的转换场景有几种:第一种,模型训练前要把特征数据切成更小的样本集,业务方说“文件太大我看不了,你转成 JSONL 我抽几行看看”;第二种,日志采集链路的上游只认 JSONL,你要把数仓里的 Parquet 结果导出给采集 agent;第三种,做数据交换,外部合作方不能用你的数据湖内核,他们只想要一个不依赖任何框架就能读的文本文件;第四种,要把数据导入 MongoDB、Elasticsearch 这类对这种“每行一个文档”格式天然友好的系统。

理解了需求之后,就能理解为什么转换不是无脑 copy 一下格式就完事。Parquet 是列式存储,把相同类型的数据扎堆存放,配合 snappy 或 zstd 压缩,压缩率可以做到文本格式的十分之一甚至更低。而 JSONL 是行式文本,每一行必须自包含,意味着原来的压缩优势基本归零。转换的过程本质上做了三件事:一是按行展开列式存储,二是把 Arrow 的类型系统映射成 JSON 的类型系统,三是重新做文本序列化。

这三个环节里,最容易出问题的是第二个。Arrow 里有 Timestamp、Decimal、List、Struct、Dictionary、二进制,还有 NaN、null、inf 这些特殊值,而 JSON 只有 null、数字、布尔、字符串、数组、对象这几样。怎么把前者安全无损地映射到后者,就是这篇文章的核心。

另外值得一提的是一次“反向直觉”的观察:很多人以为格式转换最怕的是数据量大,其实不是。数据量再大也只是时间问题,真正让任务跑挂的,往往是单行数据里的某种特殊类型。我后面会单独用一个章节讲这个,因为这块网上能查到的资料少,而且报错信息往往非常误导人。

1.1 常见工具链速览

在 Python 生态里,做这个转换的“标准套餐”其实就是两个库:pandas 和 pyarrow。性能敏感的场景还会看到 duckdb、polars、spark 的身影。

方案依赖推荐场景主要限制
pandas 的 read_parquet + to_jsonpandas, pyarrow中小文件、快速交付内存需要能装下一份半数据,嵌套类型处理麻烦
pyarrow 的 ParquetFile + iter_batchespyarrow生产环境、大文件序列化细节需要自己写
duckdb 的 COPY 语句duckdb超大文件、想少写代码嵌套类型和特殊值的行为是你不太好控制的
polars / sparkpolars / pyspark集群或超大数据量引入成本高,小需求没必要

这篇文章主要讲前两种,因为它们是绝大多数 Python 工程师最顺手的路线。duckdb 方案我会在性能优化章节提到,但不会主推——它能让你跑得很快,但遇到类型转换的边角情况时,你会发现自己很难插手。

2. 转换前的准备:环境、依赖与文件体检

这个环节经常被人跳过,但我强烈建议你不要跳。因为大多数转换失败都不是代码写错,而是你根本不了解手头的 Parquet 文件长什么样。

环境方面只需要一个 Python 3.9 以上解释器,然后装好 pandas 和 pyarrow。没有别的前置条件,不需要装 Hadoop 生态的任何东西。

pip install pandas pyarrow

装完以后,第一步永远应该是体检文件。我会写一个十几行的脚本,把文件的元信息打出来,先确认里面到底有什么类型,而不是直接闷头转换。

import pyarrow.parquet as pq import sys def inspect_parquet(path: str) -> None: pf = pq.ParquetFile(path) schema = pf.schema_arrow print(f"文件路径: {path}") print(f"总行数: {pf.metadata.num_rows}") print(f"RowGroup 数量: {pf.metadata.num_row_groups}") print(f"Schema 字段数: {len(schema.names)}") print() for field in schema: print(f" {field.name}: {field.type}") if __name__ == "__main__": inspect_parquet(sys.argv[1])

这段脚本的运行结果会直接告诉你很多东西。比如我上次遇到的一个文件,跑完发现有一列是list<item: struct<name: string, score: double>>,这种嵌套结构如果直接用 pandas 的read_parquet去读,会被 Arrow 自动压成item.nameitem.score这样的扁平列名,原始嵌套结构完全丢掉了。等到转 JSONL 的时候,输出就跟你期望的完全不一样,而且你根本不知道是哪一步改的。

2.1 为什么要做文件体检

因为 Parquet 是二进制格式,你不打开它,就永远不知道里面藏着什么。

我用“体检”这个词是有原因的。Parquet 文件里除了数据,还有一层非常重要但平常看不见的东西——Schema。它精确记录了每一列的类型、压缩方式、编码方式,甚至嵌套结构。这些信息决定了你怎么读它,也决定了你怎么写序列化逻辑。你提前知道了字段类型,后面遇到怪异的报错时就能快速判断是哪一列惹的祸。

比如跟合作方对接时,我习惯把体检脚本的输出直接贴给对方确认:“我要按这个 schema 转,细节你过目一下。”这个动作能挡掉至少一半的返工。因为实际业务里的 Parquet 文件,跟业务方口头告诉你的字段结构,经常是两码事。

2.2 关于引擎的选择:pandas 还是 pyarrow

我见过太多初学者一上来就pip install pandas,然后pd.read_parquet(...)。这不是不能用,而是你得知道自己把命运交给了谁。

pandas 的read_parquet底层有两个引擎:pyarrowfastparquet。默认情况下,pandas 会优先用 pyarrow。而 pyarrow 读取 Parquet 时,是把数据先转成 Arrow 表,再在需要的时候转成 pandas DataFrame。这个转化过程不是全免费的,尤其遇到复杂嵌套类型时,Arrow 到 pandas 的类型映射会丢失一部分原始信息。

我的建议是:如果你最终要转 JSONL,那就别绕道 pandas,直接用 pyarrow 处理,到序列化的最后一公里再碰 pandas 或者干脆不碰。因为 pandas 的to_json虽然写着方便,但它对你手头的数据类型有自己的理解,安全性和可控性都不如直接操作 Arrow 的 batch。

3. Pandas 路线实操:最直观,但你要知道它帮你扛了什么

先说这条路线为什么“直观”。因为它只需要两个方法:pd.read_parquet()df.to_json(orient='records', lines=True)。看着确实是三行搞定,但如果你文件一上来就几个 GB,这个方案很可能直接吃光你的内存。为什么?因为 pandas 在读取时会在内存里构建一份完整的 DataFrame,而to_json序列化时,即使你按行写,它其实也会先在内存里生成一个完整的 JSON 字符串列表,然后再拼接。数据量一大,这个过程的内存消耗可能是文件原始大小的好几倍。

举个例子:一个 1GB 的 Parquet 文件,压缩前数据可能有 3-4GB,读成 DataFrame 占一份,转成 JSON 字符串再占一份,加上中间过程产生的临时对象,16GB 内存的机器直接干到 swap。我见过不只一次,同事在本地跑这种脚本,电脑风扇狂转,最后进程被杀。

import pandas as pd df = pd.read_parquet("input.parquet") with open("output.jsonl", "w", encoding="utf-8") as f: f.write(df.to_json(orient="records", lines=True, force_ascii=False))

如果文件不大、内存够用、一次性的任务,这段代码足够。但我要提醒几个参数:

  • orient="records"表示每一行是一个对象,这是 JSONL 的标准形态。
  • lines=True必须配合orient="records"使用,这样输出的是每行一个 JSON 对象,而不是一整个 JSON 数组。
  • force_ascii=False让中文字符直接输出,而不是转成\uXXXX,文件体积会小很多,可读性也强得多。

但即使解决了输出格式,嵌套列你依然绕不过去。pandas 读进来的嵌套结构已经被拍平了,你再怎么调to_json的参数,也恢复不了原来的嵌套关系。所以要把这条路线真正用对,你需要给它加一个“保留嵌套”的前置条件:用 pyarrow 读取,转 pandas 时显式传types_mapper或者用to_pydict再组装,否则你会得到一个结构正确但嵌套完全丢失的结果。这就不如老老实实走 pyarrow 的自定义路线。

3.1 什么时候可以用 Pandas 路线

我给一个实用判断标准:只要你的文件在单个 GB 级别以下,schema 里没有复杂的嵌套类型,且你只是临时转一次,那 pandas 路线完全没问题。它快、简单、可读性强,尤其在本地调试、给同事发小样本的时候,效率极高。

如果超过了这个规模,或者嵌套类型很多,就往下看 PyArrow 路线。这里多说一句,实际生产环境里“单个 GB 级别”听上去小,但 Parquet 压缩率很高,1GB 的 Parquet 解压出来可能是 5-10GB 的文本。所以判断依据不是文件名大小,而是你预估解压后的 JSON 有多大,以及你机器的内存能撑到多少。

4. PyArrow 路线实操:生产级写法与完整代码

这是我最推荐的生产级方案。核心思想是:不要一次性把整个文件读进内存,而是利用ParquetFile.iter_batches()按批次读取,转一行写一行,全程内存占用保持在一个稳定的小区间。

先给完整代码,然后拆开讲为什么这么写。

import json import math from pathlib import Path from typing import Any, Callable, Iterator import pyarrow as pa import pyarrow.parquet as pq def _default_serializer(obj: Any) -> str: """处理 JSON 无法直接序列化的 Arrow/Parquet 类型。""" # 日期时间类型 if hasattr(obj, "isoformat"): return obj.isoformat() # bytes 类型,常见于 Parquet 的 binary 字段 if isinstance(obj, bytes): return obj.decode("utf-8", errors="replace") # Decimal 类型,转成字符串可以避免精度丢失 if isinstance(obj, decimal.Decimal): return str(obj) raise TypeError(f"无法序列化类型: {type(obj)}") def parquet_to_jsonl( src_path: str | Path, dst_path: str | Path, batch_size: int = 50_000, ensure_ascii: bool = False, custom_serializer: Callable[[Any], str] | None = None, ) -> int: pf = pq.ParquetFile(src_path) serializer = custom_serializer or _default_serializer total_write = 0 with open(dst_path, "w", encoding="utf-8") as out_f: for batch in pf.iter_batches(batch_size=batch_size): # 关键步骤:把 Arrow RecordBatch 转成 Python 字典列表 records = batch.to_pydict() # 将按列存储的 dict 转成按行存储的 list,再逐行 JSON 化 for i in range(batch.num_rows): row = {key: records[key][i] for key in records.keys()} line = json.dumps( row, ensure_ascii=ensure_ascii, allow_nan=False, separators=(",", ":"), default=serializer, ) out_f.write(line + "\n") total_write += 1 return total_write

这段代码看起来量不大,但里面每一个选择都有讲究。我挨个讲。

4.1 为什么用iter_batches而不是read_table

iter_batches(batch_size=...)是 PyArrow 提供的一个惰性迭代器,它不会一次性把整个文件加载进来,而是按指定的行数切块。你处理完一个 batch,这块数据就可以被垃圾回收了。内存曲线就是一条水平直线,不管底层文件是 1GB 还是 100GB,占用的内存基本恒定。

这里要注意一个容易被忽略的点:batch_size并不是越大越好。50,000 行是一个比较中庸的选择,原因有两个。一是batch.to_pydict()会一次性构造一个按列的 Python dict,行数越多,这个 dict 占的内存越大,Python 对象的开销远高于 Arrow 内部的连续内存;二是批量太大会增加单批次处理时间,一旦中途发生异常,你要损失的计算量就更大。

另外,iter_batches内部会在 RowGroup 之间自动切换,你不需要管 RowGroup 是什么。你只需要明白一个概念:Parquet 文件在物理上被切成了若干个 RowGroup,每个 RowGroup 内部的数据是连续存储的。iter_batches读取时可能会跨 RowGroup,这只是性能差异的问题,不影响输出正确性。

4.2 为什么用to_pydict()而不是to_pandas()

很多人习惯先转 pandas,再to_dict。但这里有个隐藏的坑:当 Parquet 里有嵌套结构时,batch.to_pandas()默认会把 struct 类型展开成多列,原始层级关系丢失。而batch.to_pydict()是直接按 Arrow 的内存结构输出 Python 对象,嵌套结构原样保留,不会拍平。

举个例子,如果你的 Parquet 里有一列是Struct<name: string, score: double>,那么to_pydict()会把它转成:

{"field": [{"name": "张三", "score": 92.5}, {"name": "李四", "score": 88.0}]}

而如果走了to_pandas(),你得到的可能是:

{ "field.name": ["张三", "李四"], "field.score": [92.5, 88.0] }

这两个结果对下游来说完全不是一回事。前者是一个完整的对象,后者是拍平后的列名。所以在转换 JSONL 这种天然需要结构自包含的格式时,请一定用to_pydict()

它的代价是性能。to_pydict()逐元素地把 Arrow 数组转成 Python 对象,会损失一部分 C++ 层的性能优势。但在生产场景里,这个损失的性价比是划算的,因为你换来的是类型可控、嵌套保留、零意外。

4.3json.dumps的参数都是干什么的

json.dumps有四个参数值得说清楚:

ensure_ascii=False是为了让中文直接以明文输出,而不是\u4e2d\u6587。文件体积会明显减小,查看时也更直观。但要注意,如果你的下游系统对编码有严格限制,或者你的文件要进入某些老旧的传输管道,保留默认的ensure_ascii=True可能是更安全的选择。

separators=(",", ":")是为了压缩 JSON 字符串体积。默认情况下json.dumps会在键值对之间加空格,对于千万行级别的文件,这个空格会白白多出几百 MB 的磁盘空间。压缩掉之后,人还是能一眼看懂,文件却小了一圈。

allow_nan=False是一个“强制暴露问题”的开关。Python 的json模块默认允许NaNInfinity-Infinity这种非标准 JSON 值,但你输出给下游,它们很可能无法解析。设置成False之后,一旦数据里存在 NaN 或 Infinity,脚本会立刻抛异常。这其实是好事——它帮你在转换阶段发现问题,而不是等下游解析失败再去排查。

default=serializer是一个兜底函数,处理 JSON 不认识的对象,比如datetimebytesDecimal。我在_default_serializer里写了几个最常见的分支,你也可以按自己的数据情况扩展。

4.4 关于逐行写入的性能问题

我知道有人会问:逐行out_f.write是不是太慢了?这里其实有两个层面。第一,json.dumps的耗时远大于文件写入,所以瓶颈在序列化而非写文件;第二,Python 的write调用虽然看起来是逐行,但因为有操作系统层面的缓冲,实际落盘次数远没有行数那么多。我在单机上一个 1200 万行的文件,用这段代码跑完大约 4 分钟,其中绝大部分时间花在json.dumps上。

如果你实在想优化,可以改写成批量拼接再写入:

buf = [] for i in range(batch.num_rows): ... buf.append(line) out_f.write("\n".join(buf) + "\n")

这种方式能减少 Python 层write调用的次数,对吞吐量有一点帮助。但内存占用会稍微高一点,因为你要在内存里持有整个 batch 的 JSON 字符串。建议先跑一次最朴素的版本,如果性能可以接受,就别过度优化。

5. 类型映射与序列化:转换过程中最容易踩的五个坑

前面虽然介绍了 PyArrow 路线的完整代码,但你直接拿去跑真实业务数据,大概率还是会在某个奇怪的地方挂掉。我总结了五个我自己踩过、也在同事那里反复出现的坑。每一个我都给出了“表现、原因、解决办法”三段式,方便你直接定位。

5.1 NaN、null 和 Infinity:JSON 里没有的东西

这是最经典的一个坑。Parquet 文件里的浮点列,经常存在空值或者缺失值。在 Arrow 里这些值会被表示成 null,但当你用to_pydict()转成 Python 对象时,可能会变成nanNoneinf等不同的值。json.dumpsallow_nan=False会对非标准浮点数直接抛异常,但 null 是合法的。

解决办法要看你的下游期望。如果缺失值应该输出成null,那最简单的方式是在序列化前做一次清洗:

import math def sanitize_row(row: dict) -> dict: for key, value in row.items(): if isinstance(value, float) and math.isnan(value): row[key] = None return row

注意这里math.isnan只处理浮点数,不会误伤字符串。如果值本身是整数,不需要管。如果你用的是 pandas 路线,可以考虑df = df.where(pd.notnull(df), None),但前提是你要确保 DataFrame 里的 null 不会被 pandas 转成pd.NaT

我的个人建议是:转换前先做一次全量扫描,把包含非法浮点值的行数统计出来,然后决定清洗策略。不要在代码里写死“遇到 NaN 就置 null”,因为你得先确认这些 NaN 是不是数据质量问题。

5.2 日期时间类型:时区、精度和isoformat

Arrow 的Timestamp类型在转成 Python 对象后是datetime.datetimejson.dumps默认不支持,会走default兜底。我在_default_serializer里用isoformat()解决。但这里有一个更隐蔽的问题:时间精度。

Parquet 里的时间戳精度可以是秒、毫秒、微秒、纳秒。如果你用 pandas 读取,默认会转成datetime64[ns],但isoformat()输出的是微秒精度。如果原始数据是纳秒级,你输出的时候其实已经丢精度了。这在大多数业务场景里无所谓,但如果你的下游在对比时间戳值时发现对不上,就需要回溯到这里。

另一个问题是时区。Arrow 的Timestamp可以带时区信息,转成 Pythondatetime后也可能是带时区的。JSON 里你只能输出字符串,所以建议约定一个统一的时区格式,比如都转成 UTC 的 ISO 8601 字符串。否则同一个时间,不同来源的文件可能输出2024-06-01T12:00:002024-06-01T20:00:00+08:00两种格式,下游处理起来会非常崩溃。

5.3 嵌套结构:Struct 和 List 的保真问题

前面用to_pydict()已经解决了大部分嵌套问题,但你还得面对嵌套内部的特殊类型。比如struct里套了一个timestamp,或者list里套了一个decimal。这些在 JSON 里都是合法的结构,但你需要在递归层面处理这些“叶子节点”。

给你一个通用的递归清理函数,放到代码里就能用:

def clean_value(value): if isinstance(value, dict): return {k: clean_value(v) for k, v in value.items()} if isinstance(value, (list, tuple)): return [clean_value(v) for v in value] if isinstance(value, bytes): return value.decode("utf-8", errors="replace") if hasattr(value, "isoformat"): return value.isoformat() if isinstance(value, float) and math.isnan(value): return None return value

在输出前对每个值调用clean_value(),再用json.dumps序列化,可以覆盖绝大多数类型问题。注意这个方法有两个额外好处:一是它把Decimal保留为对象时,你还需要在default兜底里处理;二是它能处理嵌套结构的递归问题,不需要你为每一层都写特殊逻辑。

5.4 字典编码列:Dictionary 转回普通值

Parquet 有一种很常见的优化手段叫字典编码——某一列的值如果重复度很高,Arrow 会把它们压缩成一个字典,数据区只存字典索引。这样读取时你会得到DictionaryArray,它的值本身是字典编码的索引,但to_pydict()会自动帮你解引用到真实值。所以这一层不需要你额外处理。

但有一个例外:当字典编码列里存在 null 值时,某些版本的 PyArrow 会把 null 和缺失值混在一起,导致你在to_pydict()结果里看到None和某个特殊标记。这时候你需要手动检查一下该列的类型:column.type是不是dictionary<...>,以及它的null_count。如果是,就明确清洗为 None。

5.5 二进制数据:Parquet 里的 bytes 怎么变成 JSON 字符串

Parquet 的Binary类型在 Python 里是bytesjson.dumps不支持 bytes,所以必须转成字符串。但转成什么字符串是有讲究的。如果你知道它是 UTF-8 编码的文本,直接decode("utf-8")即可;但如果是图片、加密串、或者别的什么,直接 decode 大概率会报错或者变成一坨乱码。

最稳妥的方法是 base64 编码,这样原数据可以无损还原:

import base64 def bytes_to_base64(value: bytes) -> str: return base64.b64encode(value).decode("ascii")

把这种函数放进clean_value或者default_serializer里,二进制字段就不会成为阻碍了。至于下游是想要可读的文本还是无损的 base64,那是业务层面的取舍,但至少你不会在这里挂掉。

6. 大数据量下的工程处理:分片、校验与断点续跑

前面讲的是单文件的正确写法。但生产环境里,我遇到的 Parquet 转换任务绝大多数不是“一个文件”,而是一个目录下几十个甚至上百个文件。这时候如果还逐文件手动跑,既不现实也容易出错。你需要一套稍微工程化的处理思路。

6.1 按文件分片输出,天然支持并行

如果一个批次的数据是多个 Parquet 文件,最简单可靠的做法是:一个输入文件对应一个输出 JSONL 文件,最后再决定要不要合并。这样做的最大好处是失败恢复极其简单——哪个文件挂了,重新跑哪个就行,不用从头再来。

import glob from pathlib import Path src_dir = Path("data/parquet") dst_dir = Path("data/jsonl") dst_dir.mkdir(parents=True, exist_ok=True) for src_file in sorted(src_dir.glob("*.parquet")): dst_file = dst_dir / f"{src_file.stem}.jsonl" if dst_file.exists(): continue # 已经处理过,跳过(断点续跑的关键) try: total = parquet_to_jsonl(src_file, dst_file) print(f"[OK] {src_file.name} -> {dst_file.name} ({total} 行)") except Exception as exc: print(f"[FAIL] {src_file.name}: {exc}")

这段代码的断点续跑能力来自if dst_file.exists(): continue。它虽然简单,但在处理几百个文件时非常管用。你可以放心地让脚本跑一半就退出,再启动时它会自动跳过已经完成的部分。

6.2 输出完成后的校验:不只是对比行数

转完之后,你拿什么证明这次转换是正确的?行数一致只是一个必要条件,不是充分条件。我建议至少做三件事:

一是检查每个输出文件的行数是否等于输入 Parquet 的num_rows。这一步用肉眼对比很痛苦,可以在脚本里自动做:

src_rows = pq.ParquetFile(src_file).metadata.num_rows dst_rows = sum(1 for _ in open(dst_file, encoding="utf-8")) assert dst_rows == src_rows, f"行数不匹配: {src_file} {src_rows} vs {dst_rows}"

二是抽检几行内容,确认字段结构符合预期。尤其是有嵌套结构的文件,转换后的结构应该跟你体检时看到的 schema 一致。抽检时优先选择文件尾部,因为尾部数据以前经常会因为增量追加而出现类型不一致的问题。

三是做一次“回读测试”:用 Python 的json.loads把输出文件的每一行都解析一遍,确认所有行都是合法 JSON。如果输出文件很大,你不可能全量回读,那就随机采样一万行做一次。这个测试能快速暴露allow_nan=False没兜住的问题。

6.3 DuckDB 路线:一杯茶的时间跑完

如果你只是想要“最快地完成转换”,不介意牺牲一点点控制权,duckdb 其实是性能怪兽。它的 SQL 语句简单到令人发指:

import duckdb duckdb.sql(""" COPY (SELECT * FROM read_parquet('data/parquet/*.parquet')) TO 'data/output.jsonl' (FORMAT json, ARRAY false) """)

就这么几行,它会自动处理分片、并行读取,速度通常比手写的 PyArrow 循环快不少。但它有两个问题:一是嵌套结构的行为跟手写方案略有差异,需要你先跑个小样本验证;二是如果数据里有非法浮点值或类型怪异的字段,它会用一种你未必预期的方式处理,而不是报错。所以在生产场景里,我通常只把 duckdb 当作“快速预览”工具,真正交付给下游还是用 PyArrow 手写方案。

6.4 输出文件要不要压缩

JSONL 转出来之后,文件体积通常比 Parquet 大 3 到 10 倍。如果你的下游不在乎明文,输出成.jsonl.gz往往比.jsonl更合适。实现起来只需要改一处:把打开输出文件的代码换成 gzip 的上下文管理器。

import gzip with gzip.open(dst_path, "wt", encoding="utf-8") as out_f: ...

需要注意,gzip 压缩是需要额外 CPU 的,而且如果下游系统不支持读 gzip,你就别用这个方案。但如果你要传输或归档,压缩率通常能到 80% 以上,非常值得。

7. 从一次真实故障说起:排查链路完整复盘

前面讲了很多理论,最后用一个真实的故障案例收尾,让你看看这些知识点在实际中是怎么连起来的。

那是一个线上跑了一个多月的定时任务,某天突然挂掉了。报错信息看起来是:

TypeError: Object of type date is not JSON serializable

第一眼看到这个错误,我心想:不就是 date 对象没处理吗?加个default函数就行。结果加完再跑,报错变成:

ValueError: Out of range float values are not JSON compliant

这就说明数据里有 NaN 或 Infinity。我用allow_nan=False把问题暴露出来后,用脚本统计了 NaN 的分布,发现是某一天新增的埋点字段里,有一列浮点数全是空值。以前的文件里这一列从来没有出现过全空的情况。

接着我去看那个字段在 Parquet 里的类型,发现它是doubleto_pydict()把空值转成了float('nan'),而之前的文件里因为至少有一个有效值,Arrow 会把它作为 null 处理,整体表现是None。全空的时候,Arrow 的字典编码逻辑就走了一条不同的分支,结果出来的就是 Python 浮点 NaN。

这个案例的核心教训是:即使同一个 Parquet 文件里的同一个字段,不同批次的数据也可能导致底层表示不同。所以转换脚本必须对“合法但不常见”的数据形态有兜底,而不是只在测试文件上验证。

排查的时候我是按这个链路走的:

  1. 先跑体检脚本,确认当前文件的 schema 和字段类型。
  2. 发现字段类型是 double,查看null_count,确认全列为空。
  3. 用一个小脚本打印to_pydict()后该列前几行的值和类型。
  4. 定位到空值被表示成nan而非None
  5. 在清洗函数里统一加math.isnan判断,输出为null

整个链路耗时不到半小时,但如果我对to_pydict()的行为没有基本预判,可能要在报错堆栈里绕很久。

这个案例也解释了我为什么反复强调:转换脚本一定要把“特殊值”当成常态来处理,而不是侥幸地认为“我们数据很干净”。数据是活的,schema 只是它某个时刻的侧写,真正流的每一批数据都可能带来惊喜。

8. 最后再分享一个实用技巧

如果你只记得这篇文章里的一件事,我希望是这件事:别让下游拿着你的 JSONL 文件去猜里面的类型。

Parquet 的关键优势之一是它自带 schema,而 JSONL 没有任何 schema 约束。同一列在第一行是字符串,第二行可能因为缺值变成了 null,下游解析时就会非常头痛。所以我在交付 JSONL 时,总会在同一目录下放一个schema.json,把每个字段的名字和类型列出来。它看起来只有几行:

{ "user_id": "string", "event_time": "timestamp", "score": "double", "tags": "array<string>" }

别小看这个文件。它帮我在多个项目里避免了下游的反复追问:你们的event_time到底是什么格式?score会不会是字符串?有了这个 schema 说明,所有疑问一次说清。

至于转换本身,我的最终建议是:先把体检脚本跑一遍,然后从 PyArrow 手写方案起步,用小样本验证输出,确认无误后再全量跑。如果你觉得每次都要写一遍太麻烦,可以把parquet_to_jsonl函数保存成一个独立模块,以后任何项目里直接 import 就能用。我自己就是把这套代码沉淀成了一个内部工具,现在每次遇到 Parquet 转 JSONL 的需求,基本都是十秒内给出答案。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/11 20:51:00

百万级QPS抢券系统架构设计与优化实践

1. 百万级QPS抢券系统的核心挑战抢券系统本质上是一个特殊的秒杀场景&#xff0c;但与普通秒杀相比存在三个显著差异点&#xff1a;首先是券的库存通常比实物商品更轻量化&#xff0c;这意味着系统可以承受更高的并发压力&#xff1b;其次是券的发放往往伴随着复杂的业务规则&a…

作者头像 李华
网站建设 2026/9/11 20:50:55

【Springboot毕设全套源码+文档】基于SpringBoot的河南文旅门户网站的设计与实现 基于SpringBoot技术的河南文旅网站开发与实现(丰富项目+远程调试+讲解+定制)

博主介绍&#xff1a;✌️码农一枚 &#xff0c;专注于大学生项目实战开发、讲解和毕业&#x1f6a2;文撰写修改等。全栈领域优质创作者&#xff0c;博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围&#xff1a;&am…

作者头像 李华
网站建设 2026/9/11 20:49:30

数值变换:修复数据与算法错位的关键一步

做数据处理这些年&#xff0c;我见过太多人拿到一批数值就直接喂进模型&#xff0c;跑出来效果不理想&#xff0c;第一反应是换算法、调参数&#xff0c;很少人回头看数据本身的数值形态。其实有相当一部分问题&#xff0c;根源不在模型强弱&#xff0c;而在输入数值的分布、量…

作者头像 李华
网站建设 2026/9/11 20:47:55

AI编程助手最新技术解析与应用实践

1. 过去一周AI Coding Agent的核心更新解析 上周AI编程助手领域迎来一波密集迭代&#xff0c;作为长期跟踪AI开发工具的技术博主&#xff0c;我梳理了各主流平台的更新亮点。不同于官方更新日志的简单罗列&#xff0c;这里将从工程实践角度分析每个改动背后的技术考量与实际影响…

作者头像 李华
网站建设 2026/9/11 20:47:53

STM32驱动AD5292数字电位器:SPI时序与寄存器配置详解

简介&#xff1a;这是一份基于STM32单片机驱动数字电位器AD5292的工程源码&#xff0c;面向嵌入式开发者和单片机爱好者&#xff0c;解决SPI接口配置与电阻值调节的实际问题。包内共2个文件&#xff0c;包含ad5292.h头文件和ad5292.c源文件&#xff0c;类型简洁&#xff0c;可直…

作者头像 李华