【问题标题】:Spark: Streaming json to parquetSpark:将json流式传输到镶木地板
【发布时间】:2016-06-16 16:18:06
【问题描述】:

如何在 Spark 流中将 json 转换为 parquet? Acutually 我必须从服务器 ssh,接收一个大的 json 文件,将其转换为镶木地板,然后将其上传到 hadoop。 我有办法以流水线的方式做到这一点吗? 它们是备份文件,所以我有一个目录,其中包含预定义数量的文件,这些文件的大小不会及时改变

类似:

scp host /dev/stdout | spark-submit myprogram.py | hadoop /dir/

编辑: 其实我正在做这个:

sc = SparkContext(appName="Test")
sqlContext = SQLContext(sc)
sqlContext.setConf("spark.sql.parquet.compression.codec.", "gzip")
#Since i couldn't get the stdio, went for a pipe:
with open("mypipe", "r") as o:
        while True:
                line = o.readline()
                print "Processing: " + line
                lineRDD = sc.parallelize([line])
                df = sqlContext.jsonRDD(lineRDD)
                #Create and append
                df.write.parquet("file:///home/user/spark/test", mode="append")
print "Done."

这工作正常,但生成的镶木地板非常大(4 行 2 列 json 为 280kb)。有什么改进吗?

【问题讨论】:

  • 只有一个文件吗?或者它是您需要从 ssh(或者可能是共享文件夹..)获取的文件流?
  • 文件有多个,但需要单独处理。
  • 是的,但它是流还是您确切知道文件的数量?这些文件是否不断被写入(作为流)到这个/这些文件夹/s?也许每隔 X 秒就会添加新文件...(?)
  • 啊,它们是备份文件,所以我有一个目录,其中包含预定义数量的文件,这些文件的大小不会及时改变。
  • 好的,很好。 (您可以更新您的问题,以便其他人也知道)我会说首先您不需要 Spark Streaming - 当然,除非您想让它更复杂。考虑到它们是预定义的,你能在运行 spark 作业之前以某种方式将这些文件复制到 HDFS 吗?也许用一个简单的脚本 fetch -> write。

标签: json apache-spark parquet


【解决方案1】:

如果有人感兴趣,我设法使用 .pipe() 方法解决了这个问题。

https://spark.apache.org/docs/latest/api/python/pyspark.html?highlight=pipe#pyspark.RDD.pipe

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-07-04
    • 2018-08-08
    • 1970-01-01
    • 1970-01-01
    • 2016-05-28
    • 2020-03-11
    • 2019-06-02
    相关资源
    最近更新 更多