【问题标题】:How to write the json file in s3 parquet如何在 s3 parquet 中写入 json 文件
【发布时间】:2020-12-19 20:13:40
【问题描述】:
    import json
    import requests
    import datetime
    import boto3
    import parquet
    import pyarrow
    import pandas as pd
    from pandas import DataFrame
         
    noaa_codes = [
        'KAST',
        'KBDN',
        'KCVO',
        'KEUG',
        'KHIO',
        'KHRI',
        'KMMV',
        'KONP',
        'KPDX',
        'KRDM',
        'KSLE',
        'KSPB',
        'KTMK',
        'KTTD',
        'KUAO'
        ]
     
    urls = [f"https://api.weather.gov/stations/{x}/observations/latest" for x in noaa_codes]
    
    
    s3_bucket="XXXXXX"
    s3_prefix = "XXXXX/parquetfiles"
    s3 = boto3.resource("s3")
    
    def get_datetime():
        dt = datetime.datetime.now()
        return dt.strftime("%Y%m%d"), dt.strftime("%H:%M:%S")
              
    def reshape(r):
        props = r["properties"]
        res = {    
            "stn": props["station"].split("/")[-1],
            "temp": props["temperature"]["value"],
            "dewp": props["dewpoint"]["value"],
            "slp": props["seaLevelPressure"]["value"],
            "stp": props["barometricPressure"]["value"],
            "visib": props["visibility"]["value"],
            "wdsp": props["windSpeed"]["value"],
            "gust": props["windGust"]["value"],
            "max": props["maxTemperatureLast24Hours"]["value"],
            "min": props["minTemperatureLast24Hours"]["value"],
            "prcp": props["precipitationLast6Hours"]["value"]
        }
        return res
               
    def lambda_handler(event, context):
                           
        responses = []
        for url in urls:
            r = requests.get(url)
            responses.append(reshape(r.json()))
        
        datestr, timestr = get_datetime()
        fname = f"noaa_hourly_measurements_{timestr}"    
        file_prefix = "/".join([s3_prefix, datestr, fname])
        s3_obj = s3.Object(s3_bucket, file_prefix)`enter code here`
        serialized = []
        for r in responses:
            serialized.append(json.dumps(r))
        jsonlines_doc = "\n".join(serialized)
        df= pd.read_json(jsonlines_doc,lines=True)
        df.to_parquet(s3_obj, engine='auto', compression='snappy', index=None)
        print("created")

无法在 aws s3 中创建 parquet 文件,但可以在本地创建。建议一个更好的方法来做到这一点。当我运行代码时,我可以在 s3 中创建一个 json 文件,但是当我尝试创建 parquet 文件时出现以下错误,出现以下错误 errorMessage": "Invalid file path or buffer object type: ", " errorType": "ValueError","stackTrace": [["/var/task/lambda_function.py",80,"lambda_handler","df.to_parquet(location, engine='auto', compression='snappy', index =无)”

【问题讨论】:

  • 我猜下面的行需要更改:df= pd.read_json(jsonlines_doc,lines=True) df.to_parquet(s3_obj, engine='auto', compression='snappy', index=无)

标签: python python-3.x pandas parquet pyarrow


【解决方案1】:

确保您的 s3_object 是一个 s3 url 字符串。它必须看起来像这样

"s3://my_bucket/path/to/data_folder/my-file.parquet"

除此之外,不建议使用 pandas 将数据帧作为 parquet 写入 S3。对于 python 3.6+,AWS 有一个名为 aws-data-wrangler 的库,它有助于 Pandas/S3/Parquet 之间的集成

安装做;

pip install awswrangler

要将您的 df 写入 s3,请执行;

import awswrangler as wr
wr.s3.to_parquet(df=df, path="s3://my_bucket/path/to/data_folder/my-file.parquet")

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-07-06
    • 1970-01-01
    • 1970-01-01
    • 2021-11-19
    • 2020-01-18
    • 2020-03-08
    • 2020-05-10
    • 2023-04-07
    相关资源
    最近更新 更多