【发布时间】:2021-06-17 06:40:04
【问题描述】:
我有一个字典列表,如下所示:
department_amount_pairs = [{"department_1": 100},{"department_2": 200},{"department_1": 300}]
我现在正在做的是
def department_udf(department_amount_pairs ):
pair = []
for d in department_amount_pairs:
pair.append(json.dumps(d))
return pair
这是我的 udf 定义
extractor = udf(department_udf,ArrayType(StringType()))
spark.udf.register("extractor_udf", extractor)
这就是我调用这个函数的方式
data = data.withColumn('pairs',extractor_udf('department_amount'))
它返回 JSON 对象.. "[{"department_1": 100},{"department_2": 200},{"department_1": 300}]" 我必须做 json.loads() 来提取这个数组。但我希望我的 udf 返回一个 Dictionaries
的 Array我尝试不使用 json.dumps 并将字典附加到列表中。但是我得到了 NONE 值..我还尝试将返回类型更改为 ArrayType(ArrayType()) 它也不起作用...
【问题讨论】:
标签: python arrays apache-spark dictionary pyspark