【问题标题】:Save and append a file in HDFS using PySpark使用 PySpark 在 HDFS 中保存和附加文件
【发布时间】:2017-11-03 00:49:54
【问题描述】:

我在 PySpark 中有一个名为 df 的数据框。我已将此df 注册为temptable,如下所示。

df.registerTempTable('mytempTable')

date=datetime.now().strftime('%Y-%m-%d %H:%M:%S')

现在我将从这个临时表中获得某些值,例如 id 列的 max_id

min_id = sqlContext.sql("select nvl(min(id),0) as minval from mytempTable").collect()[0].asDict()['minval']

max_id = sqlContext.sql("select nvl(max(id),0) as maxval from mytempTable").collect()[0].asDict()['maxval']

现在我将收集所有这些值,如下所示。

test = ("{},{},{}".format(date,min_id,max_id))

我发现test 不是data frame,而是str 字符串

>>> type(test)
<type 'str'>

现在我想将此test 保存为HDFS 中的文件。我还想将数据附加到hdfs 中的同一文件中。

如何使用 PySpark 做到这一点?

仅供参考,我使用的是 Spark 1.6,无法访问 Databricks spark-csv 包。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql hdfs


    【解决方案1】:

    在这里,您只需将您的数据与concat_ws 连接起来,并将其作为文本进行修改:

    query = """select concat_ws(',', date, nvl(min(id), 0), nvl(max(id), 0))
    from mytempTable"""
    
    sqlContext.sql(query).write("text").mode("append").save("/tmp/fooo")
    

    甚至是更好的选择:

    from pyspark.sql import functions as f
    
    (sqlContext
        .table("myTempTable")
        .select(f.concat_ws(",", f.first(f.lit(date)), f.min("id"), f.max("id")))
        .coalesce(1)
        .write.format("text").mode("append").save("/tmp/fooo"))
    

    【讨论】:

    • 这个追加到同一个目录 /tmp/fooo 是目录路径而不是文件
    • hadoop/spark 没有追加到同一个文件的东西
    • 我相信有使用hdfs dfs -appendToFile的选项
    • 你确定我们读的是同一个问题吗?这是一个恶作剧吗?
    猜你喜欢
    • 1970-01-01
    • 2019-04-23
    • 2018-11-17
    • 1970-01-01
    • 1970-01-01
    • 2021-08-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多