【发布时间】: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