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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 07:36:10