如何将show运算符的输出读回数据集?

Max*_*axU 5 scala apache-spark apache-spark-sql pyspark

假设我们有以下文本文件(df.show()命令输出):

+----+---------+--------+
|col1|     col2|    col3|
+----+---------+--------+
|   1|pi number|3.141592|
|   2| e number| 2.71828|
+----+---------+--------+
Run Code Online (Sandbox Code Playgroud)

现在我想将其解析/解析为DataFrame/Dataset.什么是最"闪亮"的方式来做到这一点?

附言:我感兴趣的解决方案 scalapyspark,这就是为什么这两个标签中使用.

Max*_*axU 4

更新:使用“UNIVOCITY”解析器库,我可以删除删除列名称中空格的一行:

斯卡拉:

// read Spark Output Fixed width table:
def readSparkOutput(filePath: String) : org.apache.spark.sql.DataFrame = {
    val t = spark.read
                 .option("header","true")
                 .option("inferSchema","true")
                 .option("delimiter","|")
                 .option("parserLib","UNIVOCITY")
                 .option("ignoreLeadingWhiteSpace","true")
                 .option("ignoreTrailingWhiteSpace","true")
                 .option("comment","+")
                 .csv(filePath)
    t.select(t.columns.filterNot(_.startsWith("_c")).map(t(_)):_*)
}
Run Code Online (Sandbox Code Playgroud)

派斯帕克:

def read_spark_output(file_path):
    t = spark.read \
             .option("header","true") \
             .option("inferSchema","true") \
             .option("delimiter","|") \
             .option("parserLib","UNIVOCITY") \
             .option("ignoreLeadingWhiteSpace","true") \
             .option("ignoreTrailingWhiteSpace","true") \
             .option("comment","+") \
             .csv("file:///tmp/spark.out")
    # select not-null columns
    return t.select([c for c in t.columns if not c.startswith("_")])
Run Code Online (Sandbox Code Playgroud)

使用示例:

scala> val df = readSparkOutput("file:///tmp/spark.out")
df: org.apache.spark.sql.DataFrame = [col1: int, col2: string ... 1 more field]

scala> df.show
+----+---------+--------+
|col1|     col2|    col3|
+----+---------+--------+
|   1|pi number|3.141592|
|   2| e number| 2.71828|
+----+---------+--------+


scala> df.printSchema
root
 |-- col1: integer (nullable = true)
 |-- col2: string (nullable = true)
 |-- col3: double (nullable = true)
Run Code Online (Sandbox Code Playgroud)

旧答案:

这是我在 scala 中的尝试(Spark 2.2):

// read Spark Output Fixed width table:
val t = spark.read
    .option("header","true")
    .option("inferSchema","true")
    .option("delimiter","|")
    .option("comment","+")
    .csv("file:///temp/spark.out")
// select not-null columns
val cols = t.columns.filterNot(c => c.startsWith("_c")).map(a => t(a))
// trim spaces from columns
val colsTrimmed = t.columns.filterNot(c => c.startsWith("_c")).map(c => c.replaceAll("\\s+",""))
// reanme columns using 'colsTrimmed'
val df = t.select(cols:_*).toDF(colsTrimmed:_*)
Run Code Online (Sandbox Code Playgroud)

它有效,但我有一种感觉,必须有更优雅的方法来做到这一点。

scala> df.show
+----+---------+--------+
|col1|     col2|    col3|
+----+---------+--------+
| 1.0|pi number|3.141592|
| 2.0| e number| 2.71828|
+----+---------+--------+

scala> df.printSchema
root
 |-- col1: double (nullable = true)
 |-- col2: string (nullable = true)
 |-- col3: double (nullable = true)
Run Code Online (Sandbox Code Playgroud)