【问题标题】:How to pass Dataframe column values dynamically from JSON file in pyspark?如何从pyspark中的JSON文件动态传递Dataframe列值?
【发布时间】:2018-09-27 08:53:18
【问题描述】:

我正在使用下面的代码创建数据框,它按预期工作。

我的数据集是'testdata'

1|123
2|223
3|323
4|423

from pyspark.sql import SQLContext,SparkSession
from pyspark.sql import Row
spark = SparkSession.builder.appName("test").getOrCreate()
sc = spark.sparkContext
sqlContext = SQLContext(sc)
df_transac = spark.createDataFrame(sc.textFile("testdata").map( lambda x: x.split("|")[:2]).map( lambda r: Row( testA = r[0],testb = r[1])))
df_transac.show()

+---------+---------+
|   testA | testB   |
+---------+---------+
|      123|        1|
|      223|        2|
|      323|        3|
|      423|        4|
+---------+---------+

以上数据框创建时间 testA,testB 是硬编码的列名,但我想从 json 中获取这些值,所以我尝试了以下方式。 我的json文件testjson.json:

{
    "column1":"testcolumn1"
    ,"column2":"testcolumn2"
}   

然后我尝试通过执行以下代码来创建数据框, 但它的抛出错误。

import json
from pyspark.sql import SQLContext,SparkSession
from pyspark.sql import Row
with open(testjson.json) as spec_data:
    jsn = json.load(spec_data)
spark = SparkSession.builder.appName("test").getOrCreate()
sc = spark.sparkContext
sqlContext = SQLContext(sc)
df_transac = spark.createDataFrame(sc.textFile("testdata").map( lambda x: x.split("|")[:2]).map( lambda r: Row( jsn['column1'] = r[0], jsn['column2'] = r[1])))

抛出错误,例如 :SyntaxError: keyword can't be an expression。

我的预期输出是:

+-----------+-----------+
|testcolumn1|testcolumn2|
+-----------+-----------+
|          1|        123|
|          2|        223|
|          3|        323|
|          4|        423|
+-----------+-----------+

请帮助我如何实现这一目标。

提前致谢。

【问题讨论】:

  • 分离出你的 createDataFrame() 调用。编写您在 createDataFrame 之外执行的所有函数,然后将 dict 对象传递给它。你会明白自己出了什么问题。

标签: python apache-spark dataframe pyspark spark-dataframe


【解决方案1】:

正如例外所说 - 您不能将表达式用作关键字,因此:

Row( jsn['column1'] = r[0], jsn['column2'] = r[1])

不是有效的 Python 代码。

您可以使用替代构造函数,然后应用参数:

Row(jsn['column1'], jsn['column2'])(r[0], r[1])

但一般情况下会更好

tmp = spark.read.option("delimiter", "|").csv("testdata")
df = tmp.select(tmp.columns[2:]).toDF(jsn['column1'], jsn['column2'])

【讨论】:

    【解决方案2】:
    df_transac = spark.createDataFrame(sc.textFile("testdata").map( lambda x: x.split("|")[:2]).map( lambda r: Row( jsn['column1'] = r[0], jsn['column2'] = r[1])))
    

    代替上面的代码,使用下面的代码来解决问题。

    c1=spec['column1']
    c2=spec['column2']
    a=sc.textFile("testdata").map( lambda x: x.split("|")[:2])
    data = sqlContext.createDataFrame(a,[c1, c2])
    

    【讨论】:

      猜你喜欢
      • 2015-08-08
      • 2018-04-24
      • 1970-01-01
      • 1970-01-01
      • 2017-01-29
      • 2020-05-04
      • 2013-06-15
      • 1970-01-01
      • 2018-07-24
      相关资源
      最近更新 更多