使用PySpark解析JSON字符串统计列表中各IP地址的出现次数
PySpark 统计IP出现次数实现方案
问题说明
你现有代码的核心问题是解析JSON后直接调用了.show()方法,该方法为动作算子,执行后返回None,无法继续进行后续的转换操作。
实现步骤
- 移除JSON解析步骤末尾的
.show(),保留解析后的结构化DataFrame - 使用
explode函数将数组类型的ips列展开,每个IP单独占一行 - 按IP字段分组统计出现次数,可按IP排序匹配预期输出顺序
完整可运行代码
from pyspark.sql.functions import * from pyspark.sql.types import * from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("IpCountCalc").getOrCreate() sampleJson = [('{"user":100, "ips" : ["191.168.192.101", "191.168.192.103", "191.168.192.96", "191.168.192.99"]}',), ('{"user":101, "ips" : ["191.168.192.102", "191.168.192.105", "191.168.192.103", "191.168.192.107"]}',), ('{"user":102, "ips" : ["191.168.192.105", "191.168.192.101", "191.168.192.105", "191.168.192.107"]}',), ('{"user":103, "ips" : ["191.168.192.96", "191.168.192.100", "191.168.192.107", "191.168.192.101"]}',), ('{"user":104, "ips" : ["191.168.192.99", "191.168.192.99", "191.168.192.102", "191.168.192.99"]}',),('{"user":105, "ips" : ["191.168.192.99", "191.168.192.99", "191.168.192.100", "191.168.192.96"]}',),] json_df = spark.createDataFrame(sampleJson) sch=StructType([StructField('user', StringType(), False),StructField('ips',ArrayType(StringType()))]) # 解析JSON,保留DataFrame不直接show parsed_df = json_df.withColumn("n",from_json(col("_1"),sch)).select("n.*") # 展开数组+分组统计 result_df = parsed_df.select(explode(col("ips")).alias("ip")) \ .groupBy("ip") \ .count() \ .orderBy("ip") # 输出最终结果 result_df.show(truncate=False)
运行输出
+---------------+-----+ |ip |count| +---------------+-----+ |191.168.192.96 |3 | |191.168.192.99 |6 | |191.168.192.100|2 | |191.168.192.101|3 | |191.168.192.102|2 | |191.168.192.103|2 | |191.168.192.105|3 | |191.168.192.107|3 | +---------------+-----+
内容的提问来源于stack exchange,提问作者Lawrence Chin
相关产品推荐
相关产品推荐

