【发布时间】:2020-01-15 01:37:58
【问题描述】:
我正在开发一个使用 FastAPI 的 API,用户可以向该 API 发出请求,以便发生以下情况:
- 首先,get 请求将从 Google Cloud Storage 获取文件并将其加载到 pyspark DataFrame 中
- 然后应用程序将对 DataFrame 执行一些转换
- 最后,我想将 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