如何用Python从GCS读取JSON.gz文件并写入PostgreSQL表
从GCS读取压缩JSON文件并写入PostgreSQL
嘿,我明白你之前处理CSV的思路,现在咱们把这套逻辑适配到JSON.gz文件上就行——核心还是下载-转换-高效写入,只是中间多了解压和JSON解析的步骤。下面是具体的实现方案,我会一步步给你讲清楚:
1. 先装必要的依赖
确保你的环境里有这些库:
pip install google-cloud-storage psycopg2-binary gzip
2. 完整代码实现
import gzip import json from io import BytesIO, StringIO from google.cloud import storage import psycopg2 import csv # 用来处理特殊字符,避免格式错误 def load_gcs_json_gz_to_postgres(): # --- 配置参数,换成你自己的 --- GCS_BUCKET = "你的存储桶名称" GCS_FILE_PATH = "path/to/your/data.json.gz" PG_CONN_INFO = { "host": "你的PG主机地址", "database": "目标数据库名", "user": "PG用户名", "password": "PG密码" } PG_TARGET_TABLE = "你的目标表名" # 表结构要和JSON字段匹配:date(TIMESTAMP), type(VARCHAR), total(NUMERIC), payment(VARCHAR) # 1. 从GCS下载压缩文件 storage_client = storage.Client() bucket = storage_client.bucket(GCS_BUCKET) blob = bucket.blob(GCS_FILE_PATH) compressed_bytes = blob.download_as_bytes() # 2. 解压并解析JSON数组 with gzip.GzipFile(fileobj=BytesIO(compressed_bytes), mode='rb') as gz_file: json_str = gz_file.read().decode('utf-8') json_data = json.loads(json_str) # 3. 生成适合COPY FROM的结构化数据流 # 用csv.writer处理特殊字符(比如字段里的制表符、引号),避免写入出错 output_stream = StringIO() csv_writer = csv.writer(output_stream, delimiter='\t', quoting=csv.QUOTE_MINIMAL) for item in json_data: # 按表字段顺序写入数据 csv_writer.writerow([ item['date'], item['type'], item['total'], item['payment'] ]) output_stream.seek(0) # 把指针移回开头,方便后续读取 # 4. 用copy_expert高效写入PostgreSQL with psycopg2.connect(**PG_CONN_INFO) as conn: with conn.cursor() as cur: copy_sql = f""" COPY {PG_TARGET_TABLE}(date, type, total, payment) FROM STDIN WITH (FORMAT csv, DELIMITER '\t', NULL ''); """ cur.copy_expert(copy_sql, output_stream) conn.commit() if __name__ == "__main__": load_gcs_json_gz_to_postgres()
3. 关键细节说明
- 解压处理:因为是gzip压缩文件,直接用
download_as_bytes()拿到字节流,再通过gzip.GzipFile解压成UTF-8字符串,才能正常解析JSON。 - 特殊字符处理:用
csv.writer生成TSV流,能自动处理字段里的换行、制表符等特殊字符,避免COPY时出现格式错误,比手动拼字符串靠谱多了。 - 大文件优化:如果你的JSON.gz文件特别大(几十GB级别),别一次性把整个JSON数组加载到内存里,推荐用
ijson库流式解析:
这种流式处理的方式不会把整个文件加载到内存,适合超大规模数据的导入。# 先装ijson:pip install ijson import ijson # 替换解压+解析的部分 with gzip.GzipFile(fileobj=BytesIO(compressed_bytes), mode='rb') as gz_file: parser = ijson.items(gz_file, 'item') # 流式遍历JSON数组的每个元素 output_stream = StringIO() csv_writer = csv.writer(output_stream, delimiter='\t') batch_size = 10000 # 每1万条数据写入一次 count = 0 with psycopg2.connect(**PG_CONN_INFO) as conn: with conn.cursor() as cur: copy_sql = f""" COPY {PG_TARGET_TABLE}(date, type, total, payment) FROM STDIN WITH (FORMAT csv, DELIMITER '\t', NULL ''); """ for item in parser: csv_writer.writerow([item['date'], item['type'], item['total'], item['payment']]) count +=1 if count % batch_size ==0: output_stream.seek(0) cur.copy_expert(copy_sql, output_stream) output_stream.truncate(0) output_stream.seek(0) # 写入剩余的最后一批数据 output_stream.seek(0) cur.copy_expert(copy_sql, output_stream) conn.commit()
4. 备选方案:直接导入JSONB字段
如果你不想拆字段,也可以把整个JSON对象写入PostgreSQL的JSONB字段,只需修改COPY逻辑:
# 生成每行一个JSON字符串的流 output_stream = StringIO() for item in json_data: output_stream.write(json.dumps(item) + '\n') output_stream.seek(0) # 修改COPY语句 copy_sql = f""" COPY {PG_TARGET_TABLE}(json_data) FROM STDIN WITH (FORMAT text); """
不过这种方式后期查询效率不如拆成独立字段高,适合临时存储或字段结构经常变化的场景。
内容的提问来源于stack exchange,提问作者emp
相关产品推荐
相关产品推荐

