使用PySpark计算Column_2中经纬度坐标列表的质心坐标
问题描述
现有数据集包含两列:
- Column_1:字符串类型,存储索引值
- Column_2:字符串数组类型,存储经纬度坐标列表(每个坐标格式为
(纬度,经度)字符串)
输入数据示例:
| Column_1 [String] | Column_2 Array[String] |
|---|---|
| A | [(40.7128,-74.0060), (34.0522, -118.2437), (51.5074, -0.1278), (-33.8651, 151.2099)] |
| B | [(52.5200, 13.4050), (48.8566, 2.3522), (41.9028, 12.4964), (55.6761, 12.5683)] |
需求:计算Column_2中每组坐标列表的质心(所有纬度的平均值、所有经度的平均值组成的坐标),输出包含原两列及Centroid列的数据集。
解决方案(Pandas版)
通过自定义函数解析坐标字符串、提取数值后计算均值,实现质心求解:
import pandas as pd # 构造输入DataFrame data = { 'Column_1': ['A', 'B'], 'Column_2': [ ['(40.7128,-74.0060)', '(34.0522, -118.2437)', '(51.5074, -0.1278)', '(-33.8651, 151.2099)'], ['(52.5200, 13.4050)', '(48.8566, 2.3522)', '(41.9028, 12.4964)', '(55.6761, 12.5683)'] ] } df = pd.DataFrame(data) # 定义质心计算函数 def calculate_centroid(coords): lats = [] lngs = [] for coord_str in coords: # 去除括号并拆分经纬度 lat_str, lng_str = coord_str.strip('()').split(',') lats.append(float(lat_str.strip())) lngs.append(float(lng_str.strip())) # 计算均值并保留4位小数 return [round(sum(lats)/len(lats), 4), round(sum(lngs)/len(lngs), 4)] # 生成Centroid列 df['Centroid'] = df['Column_2'].apply(calculate_centroid) # 输出结果 print(df.to_markdown(index=False))
运行结果
| Column_1 | Column_2 | Centroid |
|---|---|---|
| A | ['(40.7128,-74.0060)', '(34.0522, -118.2437)', '(51.5074, -0.1278)', '(-33.8651, 151.2099)'] | [48.3518, -10.0424] |
| B | ['(52.5200, 13.4050)', '(48.8566, 2.3522)', '(41.9028, 12.4964)', '(55.6761, 12.5683)'] | [49.7394, 10.4555] |
大数据场景方案(PySpark版)
如果处理大规模数据集,可使用PySpark自定义UDF实现:
from pyspark.sql import SparkSession from pyspark.sql.functions import udf from pyspark.sql.types import ArrayType, DoubleType spark = SparkSession.builder.appName("CentroidCalculation").getOrCreate() # 构造输入DataFrame data = [ ("A", ["(40.7128,-74.0060)", "(34.0522, -118.2437)", "(51.5074, -0.1278)", "(-33.8651, 151.2099)"]), ("B", ["(52.5200, 13.4050)", "(48.8566, 2.3522)", "(41.9028, 12.4964)", "(55.6761, 12.5683)"]) ] df = spark.createDataFrame(data, ["Column_1", "Column_2"]) # 定义UDF @udf(returnType=ArrayType(DoubleType())) def calculate_centroid_udf(coords): lats = [] lngs = [] for coord_str in coords: lat_str, lng_str = coord_str.strip('()').split(',') lats.append(float(lat_str.strip())) lngs.append(float(lng_str.strip())) return [sum(lats)/len(lats), sum(lngs)/len(lngs)] # 添加Centroid列 df = df.withColumn("Centroid", calculate_centroid_udf(df["Column_2"])) # 展示结果 df.show(truncate=False)
内容的提问来源于stack exchange,提问作者Sai Sumanth
相关产品推荐
相关产品推荐

