【问题标题】:Is it possible to allow users to download the result of a pyspark dataframe in FastAPI or Flask是否可以允许用户在 FastAPI 或 Flask 中下载 pyspark 数据帧的结果
【发布时间】:2020-01-15 01:37:58
【问题描述】:

我正在开发一个使用 FastAPI 的 API,用户可以向该 API 发出请求,以便发生以下情况:

  1. 首先,get 请求将从 Google Cloud Storage 获取文件并将其加载到 pyspark DataFrame 中
  2. 然后应用程序将对 DataFrame 执行一些转换
  3. 最后,我想将 DataFrame 作为 parquet 文件写入用户磁盘。

我不太清楚如何以 parquet 格式将文件传递给用户,原因如下:

  • df.write.parquet('out/path.parquet') 将数据写入out/path.parquet 的目录中,当我尝试将其传递给starlette.responses.FileResponse 时,这会带来挑战
  • 将我知道存在的单个 .parquet 文件传递​​给 starlette.responses.FileResponse 似乎只是将二进制文件打印到我的控制台(如下面的代码所示)
  • 将 DataFrame 写入 BytesIO 流 like in pandas 似乎很有希望,但我不太清楚如何使用 DataFrame 的任何方法或 DataFrame.rdd 的方法来做到这一点。

这在 FastAPI 中是否可行?是否可以在 Flask 中使用send_file()

这是我到目前为止的代码。请注意,我已经尝试了一些类似注释代码的方法,但均无济于事。

import tempfile

from fastapi import APIRouter
from pyspark.context import SparkContext
from pyspark.sql.session import SparkSession
from starlette.responses import FileResponse


router = APIRouter()
sc = SparkContext('local')
spark = SparkSession(sc)

df: spark.createDataFrame = spark.read.parquet('gs://my-bucket/sample-data/my.parquet')

@router.get("/applications")
def applications():
    df.write.parquet("temp.parquet", compression="snappy")
    return FileResponse("part-some-compressed-file.snappy.parquet")
    # with tempfile.TemporaryFile() as f:
    #     f.write(df.rdd.saveAsPickleFile("temp.parquet"))
    #     return FileResponse("test.parquet")

谢谢!

编辑: 我尝试使用here 提供的答案和信息,但我无法完全正常工作。

【问题讨论】:

    标签: python flask pyspark fastapi


    【解决方案1】:

    我能够解决这个问题,但它远非优雅。如果有人可以提供不写入磁盘的解决方案,我将不胜感激,并将选择您的答案作为正确答案。

    我能够使用df.rdd.saveAsPickleFile() 序列化DataFrame,压缩生成的目录,将其传递给python 客户端,将生成的zipfile 写入磁盘,解压缩,然后使用SparkContext().pickleFile,最后加载DataFrame。我认为远非理想。

    API:

    import shutil
    import tempfile
    
    from fastapi import APIRouter
    from pyspark.context import SparkContext
    from pyspark.sql.session import SparkSession
    from starlette.responses import FileResponse
    
    
    router = APIRouter()
    sc = SparkContext('local')
    spark = SparkSession(sc)
    
    df: spark.createDataFrame = spark.read.parquet('gs://my-bucket/my-file.parquet')
    
    @router.get("/applications")
    def applications():
        temp_parquet = tempfile.NamedTemporaryFile()
        temp_parquet.close()
        df.rdd.saveAsPickleFile(temp_parquet.name)
    
        shutil.make_archive('test', 'zip', temp_parquet.name)
    
        return FileResponse('test.zip')
    

    客户:

    import io
    import zipfile
    
    import requests
    
    from pyspark.context import SparkContext
    from pyspark.sql.session import SparkSession
    
    sc = SparkContext('local')
    spark = SparkSession(sc)
    
    response = requests.get("http://0.0.0.0:5000/applications")
    file_like_object = io.BytesIO(response.content)
    with zipfile.ZipFile(file_like_object) as z:
        z.extractall('temp.data')
    
    rdd = sc.pickleFile("temp.data")
    df = spark.createDataFrame(rdd)
    
    print(df.head())
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-09-27
      相关资源
      最近更新 更多