PySpark拆分哈希标签列时explode报错:输入需为数组/Map类型
解决PySpark中explode哈希标签列表时的类型不匹配问题
我来帮你搞定这个问题!你遇到的错误核心是UDF的返回类型没有正确声明,Spark把你的Python列表当成了字符串类型,而explode函数要求输入必须是数组(ArrayType)或映射(MapType),所以才会抛出类型不匹配的异常。
问题根源
你定义的hash_tags_fun没有指定返回类型,Spark会自动推断类型,但它把Python的列表直接序列化成了字符串(比如"[#SHINee, #AMBER]"),所以你的hash_tags_list里的text列其实是StringType,不是ArrayType,这就导致explode无法处理。
修复步骤
导入必要的类型定义
首先要引入PySpark的数组类型和字符串类型,用来声明UDF的返回值类型:from pyspark.sql.types import ArrayType, StringType重新定义带类型声明的UDF
修改你的UDF,明确指定返回类型为ArrayType(StringType()),这样Spark就知道这是一个字符串数组:hash_tags_fun = udf(lambda t: re.findall('(#[^#]\w{3,})', t), ArrayType(StringType()))执行explode操作
现在可以正常使用explode把数组中的每个哈希标签拆成单独一行了,建议给新生成的数组列起个不同的名字,避免覆盖原text列:hash_tags_result = sqlContext.sql("SELECT text FROM hash_tags_table") # 生成哈希标签数组列,命名为hash_tags hash_tags_list = hash_tags_result.withColumn('hash_tags', hash_tags_fun('text')) # 拆分数组并保留哈希标签作为text列,删除临时数组列 hash_tags_exploded = hash_tags_list.withColumn("text", explode("hash_tags")).drop("hash_tags") hash_tags_exploded.show()
执行完这些步骤后,你就能得到期望的每行一个哈希标签的结果:
+-------------------+ | text| +-------------------+ | #shutUpAndDANCE| | #SHINee| | #AMBER| | #JR50| | #flipagram| +-------------------+
验证类型(可选)
你可以用printSchema()来确认列类型是否正确:
hash_tags_list.printSchema()
修改后的输出应该显示hash_tags列是array<string>类型,而不是之前的string类型。
内容的提问来源于stack exchange,提问作者Sarath Chandra
相关产品推荐
相关产品推荐

